Files
Alessandro RosandGitHub a146e55978 webrtc: improve performance by ignoring mDNS candidates (#4963) (#6064)
mDNS candidates sometimes require a large CPU portion, they are not
involved in any connectivity method mentioned in the documentation, they work in
local networks only.
2026-08-10 12:00:23 +00:00

1087 lines
24 KiB
Go

// Package webrtc contains WebRTC utilities.
package webrtc
import (
"context"
"errors"
"fmt"
"net"
"slices"
"strconv"
"strings"
"sync"
"sync/atomic"
"time"
"github.com/pion/ice/v4"
"github.com/pion/interceptor"
"github.com/pion/sdp/v3"
"github.com/pion/transport/v4"
"github.com/pion/webrtc/v4"
"github.com/bluenviron/mediamtx/internal/logger"
)
const (
webrtcStreamID = "mediamtx"
twccExtensionURI = "http://www.ietf.org/id/draft-holmer-rmcat-transport-wide-cc-extensions-01"
)
func interfaceIPs(interfaceList []string) ([]string, error) {
intfs, err := net.Interfaces()
if err != nil {
return nil, err
}
var ips []string
for _, intf := range intfs {
if len(interfaceList) == 0 || slices.Contains(interfaceList, intf.Name) {
var addrs []net.Addr
addrs, err = intf.Addrs()
if err == nil {
for _, addr := range addrs {
var ip net.IP
switch v := addr.(type) {
case *net.IPNet:
ip = v.IP
case *net.IPAddr:
ip = v.IP
}
if ip != nil {
ips = append(ips, ip.String())
}
}
}
}
}
return ips, nil
}
func maxTrackCount(medias []*sdp.MediaDescription) int {
total := 0
for _, media := range medias {
ridCount := 0
for _, attr := range media.Attributes {
if attr.Key == "rid" {
ridCount++
}
}
if ridCount == 0 {
ridCount = 1
}
total += ridCount
}
return total
}
func ridIndexByMedia(media *sdp.MediaDescription) map[string]int {
ret := make(map[string]int)
for _, attr := range media.Attributes {
if attr.Key != "rid" || attr.Value == "" {
continue
}
ridge, _, ok := strings.Cut(attr.Value, " ")
if !ok || ridge == "" {
ridge = attr.Value
}
ret[ridge] = len(ret)
}
return ret
}
// * 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
}
err = webrtc.ConfigureSimulcastExtensionHeaders(mediaEngine)
if err != nil {
return err
}
err = webrtc.ConfigureTWCCSender(mediaEngine, interceptorRegistry)
if err != nil {
return err
}
interceptorRegistry.Add(&statsInterceptorFactory{
onCreate: onStatsInterceptor,
})
return nil
}
func candidateLabel(c *webrtc.ICECandidate) string {
return c.Typ.String() + "/" + c.Protocol.String() + "/" +
c.Address + "/" + strconv.FormatInt(int64(c.Port), 10)
}
// TracksAreValid checks whether tracks in the SDP are valid
func TracksAreValid(medias []*sdp.MediaDescription) error {
videoTrack := false
audioTrack := false
for _, media := range medias {
switch media.MediaName.Media {
case "video":
if videoTrack {
return fmt.Errorf("only a single video and a single audio track are supported")
}
videoTrack = true
case "audio":
if audioTrack {
return fmt.Errorf("only a single video and a single audio track are supported")
}
audioTrack = true
default:
return fmt.Errorf("unsupported media '%s'", media.MediaName.Media)
}
}
if !videoTrack && !audioTrack {
return fmt.Errorf("no valid tracks found")
}
return nil
}
type trackRecvPair struct {
track *webrtc.TrackRemote
receiver *webrtc.RTPReceiver
}
// PeerConnection is a wrapper around webrtc.PeerConnection.
type PeerConnection struct {
Net transport.Net
LocalRandomUDP bool
ICEUDPMux ice.UDPMux
ICETCPMux *TCPMuxWrapper
ICEServers []webrtc.ICEServer
IPsFromInterfaces bool
IPsFromInterfacesList []string
AdditionalHosts []string
STUNGatherTimeout time.Duration
SupportsIPv6 bool
Publish bool
OutboundTracks []*OutboundTrack
OutboundDataChannels []*OutboundDataChannel
Log logger.Writer
wr *webrtc.PeerConnection
ctx context.Context
ctxCancel context.CancelFunc
readingStarted atomic.Int64
inboundTracks []*InboundTrack
statsInterceptor *statsInterceptor
newLocalCandidate chan *webrtc.ICECandidateInit
inboundTrack chan trackRecvPair
stateMutex sync.Mutex
state webrtc.PeerConnectionState
stateChanged chan struct{}
gatheringMutex sync.Mutex
gatheringDone chan struct{}
done chan struct{}
chStartReading chan struct{}
}
// Start starts the peer connection.
func (co *PeerConnection) Start() error {
if co.STUNGatherTimeout == 0 {
co.STUNGatherTimeout = 5 * time.Second
}
settingsEngine := webrtc.SettingEngine{}
settingsEngine.SetIncludeLoopbackCandidate(true)
// always enable TCP since we might be the client of a remote TCP listener
networkTypes := []webrtc.NetworkType{webrtc.NetworkTypeTCP4}
if co.SupportsIPv6 {
networkTypes = append(networkTypes, webrtc.NetworkTypeTCP6)
}
if co.LocalRandomUDP || co.ICEUDPMux != nil || len(co.ICEServers) != 0 {
networkTypes = append(networkTypes, webrtc.NetworkTypeUDP4)
if co.SupportsIPv6 {
networkTypes = append(networkTypes, webrtc.NetworkTypeUDP6)
}
}
settingsEngine.SetNetworkTypes(networkTypes)
if co.ICEUDPMux != nil {
settingsEngine.SetICEUDPMux(co.ICEUDPMux)
}
if co.ICETCPMux != nil {
settingsEngine.SetICETCPMux(co.ICETCPMux.Mux)
}
settingsEngine.SetSTUNGatherTimeout(co.STUNGatherTimeout)
settingsEngine.SetNet(co.Net)
// Do not generate mDNS candidates and do not use mDNS candidates of remote peers.
// * they sometimes require a large CPU portion
// * they are not involved in any connectivity method mentioned in the documentation
// * they work in local networks only
settingsEngine.SetICEMulticastDNSMode(ice.MulticastDNSModeDisabled)
mediaEngine := &webrtc.MediaEngine{}
if co.Publish {
videoSetupped := false
audioSetupped := false
for _, tr := range co.OutboundTracks {
if tr.isVideo() {
videoSetupped = true
} else {
audioSetupped = true
}
}
// When audio is not used, a track has to be present anyway,
// otherwise video is not displayed on Firefox and Chrome.
if !audioSetupped {
co.OutboundTracks = append(co.OutboundTracks, &OutboundTrack{
Caps: webrtc.RTPCodecCapability{
MimeType: webrtc.MimeTypePCMU,
ClockRate: 8000,
},
})
}
for i, tr := range co.OutboundTracks {
var codecType webrtc.RTPCodecType
if tr.isVideo() {
codecType = webrtc.RTPCodecTypeVideo
} else {
codecType = webrtc.RTPCodecTypeAudio
}
err := mediaEngine.RegisterCodec(webrtc.RTPCodecParameters{
RTPCodecCapability: tr.Caps,
PayloadType: webrtc.PayloadType(96 + i),
}, codecType)
if err != nil {
return err
}
}
// When video is not used, a track must not be added but a codec has to present.
// Otherwise audio is muted on Firefox and Chrome.
if !videoSetupped {
err := mediaEngine.RegisterCodec(webrtc.RTPCodecParameters{
RTPCodecCapability: webrtc.RTPCodecCapability{
MimeType: webrtc.MimeTypeVP8,
ClockRate: 90000,
},
PayloadType: 96,
}, webrtc.RTPCodecTypeVideo)
if err != nil {
return err
}
}
} else {
for _, codec := range incomingVideoCodecs {
err := mediaEngine.RegisterCodec(codec, webrtc.RTPCodecTypeVideo)
if err != nil {
return err
}
}
for _, codec := range incomingAudioCodecs {
err := mediaEngine.RegisterCodec(codec, webrtc.RTPCodecTypeAudio)
if err != nil {
return err
}
}
}
interceptorRegistry := &interceptor.Registry{}
err := registerInterceptors(
mediaEngine,
interceptorRegistry,
func(s *statsInterceptor) {
co.statsInterceptor = s
},
)
if err != nil {
return err
}
api := webrtc.NewAPI(
webrtc.WithSettingEngine(settingsEngine),
webrtc.WithMediaEngine(mediaEngine),
webrtc.WithInterceptorRegistry(interceptorRegistry))
co.wr, err = api.NewPeerConnection(webrtc.Configuration{
ICEServers: co.ICEServers,
})
if err != nil {
return err
}
co.ctx, co.ctxCancel = context.WithCancel(context.Background())
co.newLocalCandidate = make(chan *webrtc.ICECandidateInit)
co.stateChanged = make(chan struct{})
co.gatheringDone = make(chan struct{})
co.inboundTrack = make(chan trackRecvPair)
co.done = make(chan struct{})
co.chStartReading = make(chan struct{})
if co.Publish {
for _, tr := range co.OutboundTracks {
err = tr.setup(co)
if err != nil {
co.wr.GracefulClose() //nolint:errcheck
return err
}
}
for _, dc := range co.OutboundDataChannels {
err = dc.setup(co)
if err != nil {
co.wr.GracefulClose() //nolint:errcheck
return err
}
}
} else {
_, err = co.wr.AddTransceiverFromKind(webrtc.RTPCodecTypeVideo, webrtc.RTPTransceiverInit{
Direction: webrtc.RTPTransceiverDirectionRecvonly,
})
if err != nil {
co.wr.GracefulClose() //nolint:errcheck
return err
}
_, err = co.wr.AddTransceiverFromKind(webrtc.RTPCodecTypeAudio, webrtc.RTPTransceiverInit{
Direction: webrtc.RTPTransceiverDirectionRecvonly,
})
if err != nil {
co.wr.GracefulClose() //nolint:errcheck
return err
}
co.wr.OnTrack(func(track *webrtc.TrackRemote, receiver *webrtc.RTPReceiver) {
select {
case co.inboundTrack <- trackRecvPair{track, receiver}:
case <-co.ctx.Done():
}
})
}
co.wr.OnConnectionStateChange(func(state webrtc.PeerConnectionState) {
co.stateMutex.Lock()
defer co.stateMutex.Unlock()
if co.state == webrtc.PeerConnectionStateFailed || co.state == webrtc.PeerConnectionStateClosed {
return
}
co.state = state
close(co.stateChanged)
co.stateChanged = make(chan struct{})
co.Log.Log(logger.Debug, "peer connection state: "+state.String())
if state == webrtc.PeerConnectionStateConnected {
co.Log.Log(logger.Info, "peer connection established, local candidate: %v, remote candidate: %v",
co.LocalCandidate(), co.RemoteCandidate())
}
})
co.wr.OnICECandidate(func(i *webrtc.ICECandidate) {
co.gatheringMutex.Lock()
defer co.gatheringMutex.Unlock()
if i != nil {
v := i.ToJSON()
select {
case co.newLocalCandidate <- &v:
case <-co.Connected():
case <-co.ctx.Done():
}
} else {
select {
case <-co.gatheringDone:
default:
close(co.gatheringDone)
}
}
})
go co.run()
return nil
}
// Close closes the connection.
func (co *PeerConnection) Close() {
co.ctxCancel()
<-co.done
}
func (co *PeerConnection) run() {
defer close(co.done)
defer func() {
for _, track := range co.inboundTracks {
track.close()
}
for _, track := range co.OutboundTracks {
track.close()
}
co.wr.GracefulClose() //nolint:errcheck
// even if GracefulClose() should wait for any goroutine to return,
// we have to wait for OnConnectionStateChange to return anyway,
// since it is executed in an uncontrolled goroutine.
// https://github.com/pion/webrtc/blob/v4.2.8/peerconnection.go#L529
<-co.failedNoContext()
}()
for {
select {
case <-co.chStartReading:
for _, track := range co.inboundTracks {
track.start()
}
co.readingStarted.Store(1)
case <-co.ctx.Done():
return
}
}
}
func (co *PeerConnection) removeUnwantedCandidates(firstMedia *sdp.MediaDescription) error {
var allowedIPs []string
if co.IPsFromInterfaces {
var err error
allowedIPs, err = interfaceIPs(co.IPsFromInterfacesList)
if err != nil {
return err
}
}
var newAttributes []sdp.Attribute //nolint:prealloc
for _, attr := range firstMedia.Attributes {
if attr.Key == "candidate" {
parts := strings.Split(attr.Value, " ")
// hide random UDP candidates
if !co.LocalRandomUDP && co.ICEUDPMux == nil && parts[2] == "udp" && parts[7] == "host" {
continue
}
// hide disallowed IPs
if parts[7] == "host" && !slices.Contains(allowedIPs, parts[4]) {
continue
}
}
newAttributes = append(newAttributes, attr)
}
firstMedia.Attributes = newAttributes
return nil
}
func (co *PeerConnection) addAdditionalCandidates(firstMedia *sdp.MediaDescription) error {
i := 0
for _, attr := range firstMedia.Attributes {
if attr.Key == "end-of-candidates" {
break
}
i++
}
for _, host := range co.AdditionalHosts {
var ips []string
if net.ParseIP(host) != nil {
ips = []string{host}
} else {
tmp, err := net.LookupIP(host)
if err != nil {
// The host can't be resolved server-side - e.g. in air-gapped
// networks without DNS, or with split-horizon / overlay DNS
// names that only resolve on the client. Skip it instead of
// failing the entire session, so the other entries still work.
co.Log.Log(logger.Warn, "cannot resolve additional host %q, skipping it: %v", host, err)
continue
}
ips = make([]string, len(tmp))
for i, e := range tmp {
ips[i] = e.String()
}
}
for _, ip := range ips {
newAttrs := append([]sdp.Attribute(nil), firstMedia.Attributes[:i]...)
if co.ICEUDPMux != nil {
port := strconv.FormatInt(int64(co.ICEUDPMux.GetListenAddresses()[0].(*net.UDPAddr).Port), 10)
tmp, err := randUint32()
if err != nil {
return err
}
id := strconv.FormatInt(int64(tmp), 10)
newAttrs = append(newAttrs, sdp.Attribute{
Key: "candidate",
Value: id + " 1 udp 2130706431 " + ip + " " + port + " typ host",
})
newAttrs = append(newAttrs, sdp.Attribute{
Key: "candidate",
Value: id + " 2 udp 2130706431 " + ip + " " + port + " typ host",
})
}
if co.ICETCPMux != nil {
port := strconv.FormatInt(int64(co.ICETCPMux.Ln.Addr().(*net.TCPAddr).Port), 10)
tmp, err := randUint32()
if err != nil {
return err
}
id := strconv.FormatInt(int64(tmp), 10)
newAttrs = append(newAttrs, sdp.Attribute{
Key: "candidate",
Value: id + " 1 tcp 1671430143 " + ip + " " + port + " typ host tcptype passive",
})
newAttrs = append(newAttrs, sdp.Attribute{
Key: "candidate",
Value: id + " 2 tcp 1671430143 " + ip + " " + port + " typ host tcptype passive",
})
}
newAttrs = append(newAttrs, firstMedia.Attributes[i:]...)
firstMedia.Attributes = newAttrs
}
}
return nil
}
func (co *PeerConnection) filterLocalDescription(desc *webrtc.SessionDescription) (*webrtc.SessionDescription, error) {
var psdp sdp.SessionDescription
psdp.Unmarshal([]byte(desc.SDP)) //nolint:errcheck
firstMedia := psdp.MediaDescriptions[0]
err := co.removeUnwantedCandidates(firstMedia)
if err != nil {
return nil, err
}
err = co.addAdditionalCandidates(firstMedia)
if err != nil {
return nil, err
}
out, _ := psdp.Marshal()
desc.SDP = string(out)
return desc, nil
}
// CreatePartialOffer creates a partial offer.
func (co *PeerConnection) CreatePartialOffer(restart bool) (*webrtc.SessionDescription, error) {
var options *webrtc.OfferOptions
if restart {
co.gatheringMutex.Lock()
select {
case <-co.gatheringDone:
default:
co.gatheringMutex.Unlock()
return nil, fmt.Errorf("tried an ICE restart before candidate gathering is complete")
}
co.gatheringDone = make(chan struct{})
co.gatheringMutex.Unlock()
options = &webrtc.OfferOptions{
ICERestart: true,
}
}
tmp, err := co.wr.CreateOffer(options)
if err != nil {
return nil, err
}
offer := &tmp
err = co.wr.SetLocalDescription(*offer)
if err != nil {
return nil, err
}
offer, err = co.filterLocalDescription(offer)
if err != nil {
return nil, err
}
return offer, nil
}
// CreateFullOffer creates a full offer.
func (co *PeerConnection) CreateFullOffer() (*webrtc.SessionDescription, error) {
tmp, err := co.wr.CreateOffer(nil)
if err != nil {
return nil, err
}
offer := &tmp
err = co.wr.SetLocalDescription(*offer)
if err != nil {
return nil, err
}
err = co.waitGatheringDone()
if err != nil {
return nil, err
}
offer = co.wr.LocalDescription()
offer, err = co.filterLocalDescription(offer)
if err != nil {
return nil, err
}
return offer, nil
}
// SetAnswer sets the answer.
func (co *PeerConnection) SetAnswer(answer *webrtc.SessionDescription) error {
return co.wr.SetRemoteDescription(*answer)
}
// RemoteDescription returns the current remote description.
func (co *PeerConnection) RemoteDescription() *webrtc.SessionDescription {
return co.wr.RemoteDescription()
}
// AddRemoteCandidate adds a remote candidate.
func (co *PeerConnection) AddRemoteCandidate(candidate *webrtc.ICECandidateInit) error {
return co.wr.AddICECandidate(*candidate)
}
// CreateFullAnswer accepts an offer and creates a full answer.
func (co *PeerConnection) CreateFullAnswer(
offer *webrtc.SessionDescription,
restarted bool,
) (*webrtc.SessionDescription, error) {
if restarted {
co.gatheringMutex.Lock()
select {
case <-co.gatheringDone:
default:
co.gatheringMutex.Unlock()
return nil, fmt.Errorf("tried an ICE restart before candidate gathering is complete")
}
co.gatheringDone = make(chan struct{})
co.gatheringMutex.Unlock()
}
err := co.wr.SetRemoteDescription(*offer)
if err != nil {
return nil, err
}
tmp, err := co.wr.CreateAnswer(nil)
if err != nil {
if errors.Is(err, webrtc.ErrSenderWithNoCodecs) {
return nil, fmt.Errorf("codecs not supported by client")
}
return nil, err
}
answer := &tmp
err = co.wr.SetLocalDescription(*answer)
if err != nil {
return nil, err
}
err = co.waitGatheringDone()
if err != nil {
return nil, err
}
answer = co.wr.LocalDescription()
answer, err = co.filterLocalDescription(answer)
if err != nil {
return nil, err
}
return answer, nil
}
func (co *PeerConnection) waitGatheringDone() error {
for {
select {
case <-co.NewLocalCandidate():
case <-co.GatheringDone():
return nil
case <-co.ctx.Done():
return fmt.Errorf("terminated")
}
}
}
// WaitUntilConnected waits until connection is established.
func (co *PeerConnection) WaitUntilConnected(timeout time.Duration) error {
t := time.NewTimer(timeout)
defer t.Stop()
outer:
for {
select {
case <-t.C:
return fmt.Errorf("deadline exceeded while waiting connection")
case <-co.Connected():
break outer
case <-co.ctx.Done():
return fmt.Errorf("terminated")
}
}
return nil
}
// GatherInboundTracks gathers incoming tracks.
func (co *PeerConnection) GatherInboundTracks(timeout time.Duration) error {
var sdp sdp.SessionDescription
sdp.Unmarshal([]byte(co.wr.RemoteDescription().SDP)) //nolint:errcheck
maxTrackCount := maxTrackCount(sdp.MediaDescriptions)
tracks := make([]*InboundTrack, 0, maxTrackCount)
t := time.NewTimer(timeout)
defer t.Stop()
midIndexByMid := make(map[string]int, len(sdp.MediaDescriptions))
ridIndexByMid := make(map[string]map[string]int, len(sdp.MediaDescriptions))
for i, media := range sdp.MediaDescriptions {
mid, _ := media.Attribute("mid")
midIndexByMid[mid] = i
ridIndexByMid[mid] = ridIndexByMedia(media)
}
getMIDIndex := func(mid string) int {
if v, ok := midIndexByMid[mid]; ok {
return v
}
return len(midIndexByMid)
}
getRIDIndex := func(mid, rid string) int {
if rid == "" {
return 0
}
if v, ok := ridIndexByMid[mid][rid]; ok {
return v
}
return len(ridIndexByMid[mid])
}
outer:
for {
select {
case <-t.C:
if len(tracks) != 0 {
break outer
}
return fmt.Errorf("deadline exceeded while waiting tracks")
case pair := <-co.inboundTrack:
mid := ""
if transceiver := pair.receiver.RTPTransceiver(); transceiver != nil {
mid = transceiver.Mid()
}
rid := pair.track.RID()
track := &InboundTrack{
track: pair.track,
receiver: pair.receiver,
midIndex: getMIDIndex(mid),
rid: rid,
ridIndex: getRIDIndex(mid, rid),
writeRTCP: co.wr.WriteRTCP,
log: co.Log,
}
track.initialize()
tracks = append(tracks, track)
if len(tracks) >= maxTrackCount {
break outer
}
case <-co.Failed():
return fmt.Errorf("peer connection closed")
case <-co.ctx.Done():
return fmt.Errorf("terminated")
}
}
slices.SortStableFunc(tracks, func(track1 *InboundTrack, track2 *InboundTrack) int {
if track1.midIndex != track2.midIndex {
return track1.midIndex - track2.midIndex
}
if track1.ridIndex != track2.ridIndex {
return track1.ridIndex - track2.ridIndex
}
id1 := track1.track.ID()
id2 := track2.track.ID()
if id1 != id2 {
return strings.Compare(id1, id2)
}
streamID1 := track1.track.StreamID()
streamID2 := track2.track.StreamID()
return strings.Compare(streamID1, streamID2)
})
co.inboundTracks = tracks
return nil
}
// Connected returns when connected.
func (co *PeerConnection) Connected() <-chan struct{} {
ch := make(chan struct{})
go func() {
for {
co.stateMutex.Lock()
state := co.state
stateChanged := co.stateChanged
co.stateMutex.Unlock()
if state == webrtc.PeerConnectionStateConnected {
close(ch)
return
}
select {
case <-stateChanged:
case <-co.ctx.Done():
// exit without closing ch
return
}
}
}()
return ch
}
// Failed returns when failed or closed.
func (co *PeerConnection) Failed() <-chan struct{} {
ch := make(chan struct{})
go func() {
defer close(ch)
for {
co.stateMutex.Lock()
state := co.state
stateChanged := co.stateChanged
co.stateMutex.Unlock()
// "closed" can arrive before "failed" and without
// the Close() method being called at all.
// It happens when the other peer sends a termination
// message like a DTLS CloseNotify.
if state == webrtc.PeerConnectionStateFailed || state == webrtc.PeerConnectionStateClosed {
return
}
select {
case <-stateChanged:
case <-co.ctx.Done():
return
}
}
}()
return ch
}
func (co *PeerConnection) failedNoContext() <-chan struct{} {
ch := make(chan struct{})
go func() {
for {
co.stateMutex.Lock()
state := co.state
stateChanged := co.stateChanged
co.stateMutex.Unlock()
if state == webrtc.PeerConnectionStateFailed || state == webrtc.PeerConnectionStateClosed {
close(ch)
return
}
<-stateChanged
}
}()
return ch
}
// NewLocalCandidate returns when there's a new local candidate.
func (co *PeerConnection) NewLocalCandidate() <-chan *webrtc.ICECandidateInit {
return co.newLocalCandidate
}
// GatheringDone returns when candidate gathering is complete.
func (co *PeerConnection) GatheringDone() <-chan struct{} {
return co.gatheringDone
}
// InboundTracks returns incoming tracks.
func (co *PeerConnection) InboundTracks() []*InboundTrack {
return co.inboundTracks
}
// StartReading starts reading incoming tracks.
func (co *PeerConnection) StartReading() {
select {
case co.chStartReading <- struct{}{}:
case <-co.ctx.Done():
}
}
// LocalCandidate returns the local candidate.
func (co *PeerConnection) LocalCandidate() string {
receivers := co.wr.GetReceivers()
if len(receivers) < 1 {
return ""
}
cp, err := receivers[0].Transport().ICETransport().GetSelectedCandidatePair()
if err != nil || cp == nil {
return ""
}
return candidateLabel(cp.Local)
}
// RemoteCandidate returns the remote candidate.
func (co *PeerConnection) RemoteCandidate() string {
receivers := co.wr.GetReceivers()
if len(receivers) < 1 {
return ""
}
cp, err := receivers[0].Transport().ICETransport().GetSelectedCandidatePair()
if err != nil || cp == nil {
return ""
}
return candidateLabel(cp.Remote)
}
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, tstats.BytesSent
}
}
}
return 0, 0
}
// Stats returns statistics.
func (co *PeerConnection) Stats() *Stats {
bytesReceived, bytesSent := bytesStats(co.wr)
v := float64(0)
n := float64(0)
packetsReceived := uint64(0)
packetsSent := uint64(0)
packetsLost := uint64(0)
if co.readingStarted.Load() == 1 {
for _, tr := range co.inboundTracks {
if recvStats := tr.rtpReceiver.Stats(); recvStats != nil {
v += recvStats.Jitter
n++
packetsReceived += recvStats.Received
packetsLost += recvStats.Lost
}
}
}
for _, tr := range co.OutboundTracks {
if sentStats := tr.rtcpSender.Stats(); sentStats != nil {
packetsSent += sentStats.Sent
}
}
var rtpPacketsJitter float64
if n != 0 {
rtpPacketsJitter = v / n
} else {
rtpPacketsJitter = 0
}
return &Stats{
BytesReceived: bytesReceived,
BytesSent: bytesSent,
RTPPacketsReceived: packetsReceived,
RTPPacketsSent: packetsSent,
RTPPacketsLost: packetsLost,
RTPPacketsJitter: rtpPacketsJitter,
RTCPPacketsReceived: co.statsInterceptor.rtcpPacketsReceived.Load(),
RTCPPacketsSent: co.statsInterceptor.rtcpPacketsSent.Load(),
}
}