diff --git a/CHANGELOG.md b/CHANGELOG.md index 864055f60..aa86675f3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,48 @@ This project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [1.13.7] - 2026-09-14 + +### Added + +- Participant kind agent details (#4809) +- report a reason on room-ended telemetry (#4815) +- telemetry: support roomID change for a participant (#4816) +- log telemetry guard details (#4817) +- ingress: add opt-in support for udp:// URL pull ingress (#4810) +- Support VP9/AV1 simulcast (#4830) +- add test-server support for compressed join requests (#4834) +- wire psrpc bus compression into the message bus constructor (#4844) +- turn: accept PROXY protocol on the TCP listener (#4852) +- telemetry: add visibility into stats worker reference underflow (#4856) +- rtc: add GetSubscribedDataTracks to LocalParticipant (#4862) + +### Changed + +- Skip the docker-backed service tests when there is no docker (#4799) +- Wait for the callbacks these tests assert on. (#4803) +- utils: make Median generic, overflow-safe, and add tests (#4553) +- warp log (#4811) +- Remove duplication within cpuload and sysload (#4818) +- Update github.com/livekit/protocol digest to a469dd4 (#4820) +- update ice fork for warp (#4824) +- Move data track packet serialization to protocol (#4829) +- Update README to include snippet about Agents (#4808) +- protocol update for SDP unmarshal hardening (#4836) +- Update golang Docker tag to v1.26.7 (#4846) +- README rewrite (#4837) + +### Fixed + +- Set up track info properly for dummy receiver. (#4807) +- Do not escape the ICE server URI in the WHIP Link header (#4795) +- fix: hold signal messages until the ReconnectResponse goes out (#4827) +- Fix/reconcile data track subscription deadlock (#4843) +- telemetry: drop unused getCPUStats, unblocking darwin builds without cgo (#4841) +- update psrpc for subscription close goroutine leak fix (#4859) +- telemetry: do not recreate a stats worker for a released guard (#4860) +- Get participant by authed identity in WHIP participant service. (#4861) + ## [1.13.6] - 2026-08-26 ### Added diff --git a/go.mod b/go.mod index db544e944..3abe59221 100644 --- a/go.mod +++ b/go.mod @@ -36,7 +36,7 @@ require ( github.com/pion/rtcp v1.2.17 github.com/pion/rtp v1.10.5 github.com/pion/sctp v1.11.1 - github.com/pion/sdp/v3 v3.0.19 + github.com/pion/sdp/v3 v3.0.20 github.com/pion/transport/v4 v4.1.0 github.com/pion/turn/v5 v5.0.13 github.com/pion/webrtc/v4 v4.2.18 @@ -77,7 +77,7 @@ require ( github.com/go-logr/stdr v1.2.2 // indirect github.com/goccy/go-json v0.10.6 // indirect github.com/golang-jwt/jwt/v5 v5.3.1 // indirect - github.com/gotesttools/gotestfmt/v2 v2.4.1 // indirect + github.com/gotesttools/gotestfmt/v2 v2.5.0 // indirect github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 // indirect github.com/mattn/go-colorable v0.1.15 // indirect github.com/mattn/go-isatty v0.0.22 // indirect @@ -147,7 +147,7 @@ require ( github.com/prometheus/client_model v0.6.3 // indirect github.com/prometheus/common v0.71.0 // indirect github.com/prometheus/procfs v0.22.0 // indirect - github.com/urfave/cli/v3 v3.10.1 + github.com/urfave/cli/v3 v3.11.0 github.com/wlynxg/anet v0.0.5 // indirect github.com/zeebo/xxh3 v1.1.0 // indirect go.uber.org/zap/exp v0.3.0 // indirect diff --git a/go.sum b/go.sum index b787420fd..56204eb47 100644 --- a/go.sum +++ b/go.sum @@ -107,8 +107,8 @@ github.com/google/wire v0.7.0 h1:JxUKI6+CVBgCO2WToKy/nQk0sS+amI9z9EjVmdaocj4= github.com/google/wire v0.7.0/go.mod h1:n6YbUQD9cPKTnHXEBN2DXlOp/mVADhVErcMFb0v3J18= github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg= github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE= -github.com/gotesttools/gotestfmt/v2 v2.4.1 h1:Ml+KPqPocp/KckpizL+tgsy/dlddI4/z2w6lgS7YIFE= -github.com/gotesttools/gotestfmt/v2 v2.4.1/go.mod h1:oQJg2KZ2aGoqEbMC2PDaAeBYm0tOkocgixK9FzsCdp4= +github.com/gotesttools/gotestfmt/v2 v2.5.0 h1:fSU3MnR+E+fvuXdw1l8xbufKhDxY3Tfjsjx/I1WerB4= +github.com/gotesttools/gotestfmt/v2 v2.5.0/go.mod h1:oQJg2KZ2aGoqEbMC2PDaAeBYm0tOkocgixK9FzsCdp4= github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 h1:5VipnvEpbqr2gA2VbM+nYVbkIF28c5ZQfqCBQ5g2xfk= github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0/go.mod h1:Hyl3n6Twe1hvtd9XUXDec4pTvgMSEixRuQKPTMH2bNs= github.com/hashicorp/go-cleanhttp v0.5.2 h1:035FKYIWjmULyFRBKPs8TBQoi0x6d9G4xc9neXJWAZQ= @@ -257,8 +257,8 @@ github.com/pion/rtp v1.10.5 h1:ip0HhO/wYZqQ4bKS+R99KnZh/GRCmIT0jDXikub7vlE= github.com/pion/rtp v1.10.5/go.mod h1:Au8fc6cEByy8RLTwKTQTEeQqDB/SJDxwL4mZuxYA5Pk= github.com/pion/sctp v1.11.1 h1:O4dIFyURw1KTST7w+gtD4gLeYXkhPa0xXLHMMoe/OSA= github.com/pion/sctp v1.11.1/go.mod h1:7KFmTwLcoYgJs/Z+99nJvsWL0qDpuyloSI0RbAqlrz0= -github.com/pion/sdp/v3 v3.0.19 h1:1VMKs3gIkTQV5M3hNKfTAPrDXSNrYtOlmOD8+mSZUGQ= -github.com/pion/sdp/v3 v3.0.19/go.mod h1:dE5WOSlzXrtiE/iuZqe9n+AcEbOjtAd3k5m5NtlV/qU= +github.com/pion/sdp/v3 v3.0.20 h1:TS6DViqcmp+49f0+mjw9anbr9xY3vJtsZewxAvlMCRQ= +github.com/pion/sdp/v3 v3.0.20/go.mod h1:slIMXDK5OKj0nhISwjfeN18AzTBCt2LYZq9uPw0cU5Q= github.com/pion/srtp/v3 v3.0.13 h1:FmQaqgNbN1vUtMhEsmj8trldc3lNZr1xmN7nl8CyX+Q= github.com/pion/srtp/v3 v3.0.13/go.mod h1:7qR3L69t8RX0EPVQwGNwCa1Gy9keKKNDpWwQzZbeXDY= github.com/pion/stun/v3 v3.1.7 h1:uRXMTlGLf89WgItGNyZ6aR5jMTX0NBbybXADpQCzn+E= @@ -319,8 +319,8 @@ github.com/twitchtv/twirp v8.1.3+incompatible h1:+F4TdErPgSUbMZMwp13Q/KgDVuI7HJX github.com/twitchtv/twirp v8.1.3+incompatible/go.mod h1:RRJoFSAmTEh2weEqWtpPE3vFK5YBhA6bqp2l1kfCC5A= github.com/ua-parser/uap-go v0.0.0-20260529044130-17c35e68e58c h1:XbG4n3OWA1PcRTpbBA22E2ChPLvJCuwYRXO12tIyVL0= github.com/ua-parser/uap-go v0.0.0-20260529044130-17c35e68e58c/go.mod h1:gwANdYmo9R8LLwGnyDFWK2PMsaXXX2HhAvCnb/UhZsM= -github.com/urfave/cli/v3 v3.10.1 h1:7Kx9H50hrHbRbyxgO1KP6/BcbiGRz0uYh5YyQ30JEEY= -github.com/urfave/cli/v3 v3.10.1/go.mod h1:ysVLtOEmg2tOy6PknnYVhDoouyC/6N42TMeoMzskhso= +github.com/urfave/cli/v3 v3.11.0 h1:P/euJp99kb9p0tlVY+iYTLYYTAQlfl0hR2gUO1Img1Q= +github.com/urfave/cli/v3 v3.11.0/go.mod h1:ysVLtOEmg2tOy6PknnYVhDoouyC/6N42TMeoMzskhso= github.com/urfave/negroni/v3 v3.1.1 h1:6MS4nG9Jk/UuCACaUlNXCbiKa0ywF9LXz5dGu09v8hw= github.com/urfave/negroni/v3 v3.1.1/go.mod h1:jWvnX03kcSjDBl/ShB0iHvx5uOs7mAzZXW+JvJ5XYAs= github.com/wlynxg/anet v0.0.5 h1:J3VJGi1gvo0JwZ/P1/Yc/8p63SoW98B5dHkYDmpgvvU= diff --git a/pkg/sfu/codecmunger/codecmunger.go b/pkg/sfu/codecmunger/codecmunger.go index 413af9688..47566a96e 100644 --- a/pkg/sfu/codecmunger/codecmunger.go +++ b/pkg/sfu/codecmunger/codecmunger.go @@ -26,6 +26,9 @@ var ( ErrFilteredVP8TemporalLayer = errors.New("filtered VP8 temporal layer") ) +// MaxHeaderSize is the largest codec header the mungers produce (VP8: 6 bytes). +const MaxHeaderSize = 8 + type CodecMunger interface { GetState() any SeedState(state any) @@ -33,7 +36,9 @@ type CodecMunger interface { SetLast(extPkt *buffer.ExtPacket) UpdateOffsets(extPkt *buffer.ExtPacket) - UpdateAndGet(extPkt *buffer.ExtPacket, snOutOfOrder bool, snHasGap bool, maxTemporal int32) (int, []byte, error) + // UpdateAndGet returns the incoming codec header size and the munged + // header by value, so the per-packet path does not allocate. + UpdateAndGet(extPkt *buffer.ExtPacket, snOutOfOrder bool, snHasGap bool, maxTemporal int32) (int, [MaxHeaderSize]byte, int, error) UpdateAndGetPadding(newPicture bool) ([]byte, error) } diff --git a/pkg/sfu/codecmunger/null.go b/pkg/sfu/codecmunger/null.go index 616786a69..41f5cb405 100644 --- a/pkg/sfu/codecmunger/null.go +++ b/pkg/sfu/codecmunger/null.go @@ -45,8 +45,8 @@ func (n *Null) SetLast(_extPkt *buffer.ExtPacket) { func (n *Null) UpdateOffsets(_extPkt *buffer.ExtPacket) { } -func (n *Null) UpdateAndGet(_extPkt *buffer.ExtPacket, snOutOfOrder bool, snHasGap bool, maxTemporal int32) (int, []byte, error) { - return 0, nil, nil +func (n *Null) UpdateAndGet(_extPkt *buffer.ExtPacket, snOutOfOrder bool, snHasGap bool, maxTemporal int32) (int, [MaxHeaderSize]byte, int, error) { + return 0, [MaxHeaderSize]byte{}, 0, nil } func (n *Null) UpdateAndGetPadding(newPicture bool) ([]byte, error) { diff --git a/pkg/sfu/codecmunger/vp8.go b/pkg/sfu/codecmunger/vp8.go index f6f45e1e6..4ffc7e65a 100644 --- a/pkg/sfu/codecmunger/vp8.go +++ b/pkg/sfu/codecmunger/vp8.go @@ -153,10 +153,11 @@ func (v *VP8) UpdateOffsets(extPkt *buffer.ExtPacket) { v.exemptedPictureIds = orderedmap.NewOrderedMap[int32, bool]() } -func (v *VP8) UpdateAndGet(extPkt *buffer.ExtPacket, snOutOfOrder bool, snHasGap bool, maxTemporalLayer int32) (int, []byte, error) { +func (v *VP8) UpdateAndGet(extPkt *buffer.ExtPacket, snOutOfOrder bool, snHasGap bool, maxTemporalLayer int32) (int, [MaxHeaderSize]byte, int, error) { + var hdr [MaxHeaderSize]byte vp8, ok := extPkt.Payload.(codec.VP8) if !ok { - return 0, nil, ErrNotVP8 + return 0, hdr, 0, ErrNotVP8 } extPictureId := v.pictureIdWrapHandler.Unwrap(vp8.PictureID, vp8.M) @@ -165,7 +166,7 @@ func (v *VP8) UpdateAndGet(extPkt *buffer.ExtPacket, snOutOfOrder bool, snHasGap if snOutOfOrder { pictureIdOffset, ok := v.missingPictureIds.Get(extPictureId) if !ok { - return 0, nil, ErrOutOfOrderVP8PictureIdCacheMiss + return 0, hdr, 0, ErrOutOfOrderVP8PictureIdCacheMiss } // the out-of-order picture id cannot be deleted from the cache @@ -190,11 +191,11 @@ func (v *VP8) UpdateAndGet(extPkt *buffer.ExtPacket, snOutOfOrder bool, snHasGap IsKeyFrame: vp8.IsKeyFrame, HeaderSize: vp8.HeaderSize + codec.VPxPictureIdSizeDiff(mungedPictureId > 127, vp8.M), } - vp8HeaderBytes, err := vp8Packet.Marshal() + n, err := vp8Packet.MarshalTo(hdr[:]) if err != nil { - return 0, nil, err + return 0, hdr, 0, err } - return vp8.HeaderSize, vp8HeaderBytes, nil + return vp8.HeaderSize, hdr, n, nil } prevMaxPictureId := v.pictureIdWrapHandler.MaxPictureId() @@ -262,7 +263,7 @@ func (v *VP8) UpdateAndGet(extPkt *buffer.ExtPacket, snOutOfOrder bool, snHasGap v.pictureIdOffset += 1 } - return 0, nil, ErrFilteredVP8TemporalLayer + return 0, hdr, 0, ErrFilteredVP8TemporalLayer } } } @@ -297,11 +298,11 @@ func (v *VP8) UpdateAndGet(extPkt *buffer.ExtPacket, snOutOfOrder bool, snHasGap IsKeyFrame: vp8.IsKeyFrame, HeaderSize: vp8.HeaderSize + codec.VPxPictureIdSizeDiff(mungedPictureId > 127, vp8.M), } - vp8HeaderBytes, err := vp8Packet.Marshal() + n, err := vp8Packet.MarshalTo(hdr[:]) if err != nil { - return 0, nil, err + return 0, hdr, 0, err } - return vp8.HeaderSize, vp8HeaderBytes, nil + return vp8.HeaderSize, hdr, n, nil } func (v *VP8) UpdateAndGetPadding(newPicture bool) ([]byte, error) { diff --git a/pkg/sfu/codecmunger/vp8_test.go b/pkg/sfu/codecmunger/vp8_test.go index 2c5714bed..7c6d1a36f 100644 --- a/pkg/sfu/codecmunger/vp8_test.go +++ b/pkg/sfu/codecmunger/vp8_test.go @@ -195,11 +195,11 @@ func TestOutOfOrderPictureId(t *testing.T) { vp8.PictureID = 13466 extPkt, _ = testutils.GetTestExtPacketVP8(params, vp8) - nIn, buf, err := v.UpdateAndGet(extPkt, true, false, 2) + nIn, hdr, nOut, err := v.UpdateAndGet(extPkt, true, false, 2) require.Error(t, err) require.ErrorIs(t, err, ErrOutOfOrderVP8PictureIdCacheMiss) require.Equal(t, 0, nIn) - require.Nil(t, buf) + require.Equal(t, 0, nOut) // create a hole in picture id vp8.PictureID = 13469 @@ -222,10 +222,10 @@ func TestOutOfOrderPictureId(t *testing.T) { } marshalledVP8, err := expectedVP8.Marshal() require.NoError(t, err) - nIn, buf, err = v.UpdateAndGet(extPkt, false, true, 2) + nIn, hdr, nOut, err = v.UpdateAndGet(extPkt, false, true, 2) require.NoError(t, err) require.Equal(t, 6, nIn) - require.Equal(t, marshalledVP8, buf) + require.Equal(t, marshalledVP8, hdr[:nOut]) // all three, the last, the current and the in-between should have been added to missing picture id cache value, ok := v.PictureIdOffset(13467) @@ -261,10 +261,10 @@ func TestOutOfOrderPictureId(t *testing.T) { } marshalledVP8, err = expectedVP8.Marshal() require.NoError(t, err) - nIn, buf, err = v.UpdateAndGet(extPkt, true, false, 2) + nIn, hdr, nOut, err = v.UpdateAndGet(extPkt, true, false, 2) require.NoError(t, err) require.Equal(t, 6, nIn) - require.Equal(t, marshalledVP8, buf) + require.Equal(t, marshalledVP8, hdr[:nOut]) } func TestTemporalLayerFiltering(t *testing.T) { @@ -294,11 +294,11 @@ func TestTemporalLayerFiltering(t *testing.T) { v.SetLast(extPkt) // translate - nIn, buf, err := v.UpdateAndGet(extPkt, false, false, 0) + nIn, _, nOut, err := v.UpdateAndGet(extPkt, false, false, 0) require.Error(t, err) require.ErrorIs(t, err, ErrFilteredVP8TemporalLayer) require.Equal(t, 0, nIn) - require.Nil(t, buf) + require.Equal(t, 0, nOut) dropped, _ := v.droppedPictureIds.Get(13467) require.True(t, dropped) require.EqualValues(t, 1, v.pictureIdOffset) @@ -308,11 +308,11 @@ func TestTemporalLayerFiltering(t *testing.T) { params.SequenceNumber = 23334 extPkt, _ = testutils.GetTestExtPacketVP8(params, vp8) - nIn, buf, err = v.UpdateAndGet(extPkt, false, false, 0) + nIn, _, nOut, err = v.UpdateAndGet(extPkt, false, false, 0) require.Error(t, err) require.ErrorIs(t, err, ErrFilteredVP8TemporalLayer) require.Equal(t, 0, nIn) - require.Nil(t, buf) + require.Equal(t, 0, nOut) dropped, _ = v.droppedPictureIds.Get(13467) require.True(t, dropped) require.EqualValues(t, 1, v.pictureIdOffset) @@ -322,11 +322,11 @@ func TestTemporalLayerFiltering(t *testing.T) { params.SequenceNumber = 23337 extPkt, _ = testutils.GetTestExtPacketVP8(params, vp8) - nIn, buf, err = v.UpdateAndGet(extPkt, false, false, 0) + nIn, _, nOut, err = v.UpdateAndGet(extPkt, false, false, 0) require.Error(t, err) require.ErrorIs(t, err, ErrFilteredVP8TemporalLayer) require.Equal(t, 0, nIn) - require.Nil(t, buf) + require.Equal(t, 0, nOut) dropped, _ = v.droppedPictureIds.Get(13467) require.True(t, dropped) require.EqualValues(t, 1, v.pictureIdOffset) @@ -376,10 +376,10 @@ func TestGapInSequenceNumberSamePicture(t *testing.T) { } marshalledVP8, err := expectedVP8.Marshal() require.NoError(t, err) - nIn, buf, err := v.UpdateAndGet(extPkt, false, false, 2) + nIn, hdr, nOut, err := v.UpdateAndGet(extPkt, false, false, 2) require.NoError(t, err) require.Equal(t, 6, nIn) - require.Equal(t, marshalledVP8, buf) + require.Equal(t, marshalledVP8, hdr[:nOut]) // telling there is a gap in sequence number will add pictures to missing picture cache expectedVP8 = &codec.VP8{ @@ -399,10 +399,10 @@ func TestGapInSequenceNumberSamePicture(t *testing.T) { } marshalledVP8, err = expectedVP8.Marshal() require.NoError(t, err) - nIn, buf, err = v.UpdateAndGet(extPkt, false, true, 2) + nIn, hdr, nOut, err = v.UpdateAndGet(extPkt, false, true, 2) require.NoError(t, err) require.Equal(t, 6, nIn) - require.Equal(t, marshalledVP8, buf) + require.Equal(t, marshalledVP8, hdr[:nOut]) value, ok := v.PictureIdOffset(13467) require.True(t, ok) diff --git a/pkg/sfu/downtrack.go b/pkg/sfu/downtrack.go index 1b993158c..d83c020f4 100644 --- a/pkg/sfu/downtrack.go +++ b/pkg/sfu/downtrack.go @@ -1044,10 +1044,11 @@ func (d *DownTrack) WriteRTP(extPkt *buffer.ExtPacket, layer int32) int32 { return 0 } + codecBytes := tp.codecHeader() poolEntity := PacketFactory.Get().(*[]byte) payload := *poolEntity - copy(payload, tp.codecBytes) - n := copy(payload[len(tp.codecBytes):], extPkt.Packet.Payload[tp.incomingHeaderSize:]) + copy(payload, codecBytes) + n := copy(payload[len(codecBytes):], extPkt.Packet.Payload[tp.incomingHeaderSize:]) if n != len(extPkt.Packet.Payload[tp.incomingHeaderSize:]) { d.params.Logger.Errorw( "payload overflow", errPayloadOverflow, @@ -1057,7 +1058,7 @@ func (d *DownTrack) WriteRTP(extPkt *buffer.ExtPacket, layer int32) int32 { PacketFactory.Put(poolEntity) return 0 } - payload = payload[:len(tp.codecBytes)+n] + payload = payload[:len(codecBytes)+n] trailerStripped := 0 if d.params.StripPacketTrailer { @@ -1068,7 +1069,7 @@ func (d *DownTrack) WriteRTP(extPkt *buffer.ExtPacket, layer int32) int32 { // translate RTP header hdr := RTPHeaderFactory.Get().(*rtp.Header) - *hdr = rtp.Header{ + initPooledRTPHeader(hdr, rtp.Header{ Version: extPkt.Packet.Version, Padding: false, Marker: tp.marker, @@ -1076,7 +1077,7 @@ func (d *DownTrack) WriteRTP(extPkt *buffer.ExtPacket, layer int32) int32 { SequenceNumber: uint16(tp.rtp.extSequenceNumber), Timestamp: uint32(tp.rtp.extTimestamp), SSRC: d.ssrc, - } + }) // add extensions if d.dependencyDescriptorExtID != 0 && tp.ddBytes != nil { @@ -1129,7 +1130,7 @@ func (d *DownTrack) WriteRTP(extPkt *buffer.ExtPacket, layer int32) int32 { tp.rtp.extTimestamp, hdr.Marker, int8(layer), - payload[:len(tp.codecBytes)], + payload[:len(codecBytes)], tp.incomingHeaderSize, tp.ddBytes, actBytes, @@ -1280,7 +1281,7 @@ func (d *DownTrack) WritePaddingRTP(bytesToSend int, paddingOnMute bool, forceMa payloads := make([]byte, RTPPaddingMaxPayloadSize*len(snts)) for i := range snts { hdr := RTPHeaderFactory.Get().(*rtp.Header) - *hdr = rtp.Header{ + initPooledRTPHeader(hdr, rtp.Header{ Version: 2, Padding: true, PaddingSize: byte(RTPPaddingMaxPayloadSize), @@ -1289,7 +1290,7 @@ func (d *DownTrack) WritePaddingRTP(bytesToSend int, paddingOnMute bool, forceMa SequenceNumber: uint16(snts[i].extSequenceNumber), Timestamp: uint32(snts[i].extTimestamp), SSRC: d.ssrc, - } + }) d.addDummyExtensions(hdr) payload := payloads[i*RTPPaddingMaxPayloadSize : (i+1)*RTPPaddingMaxPayloadSize : (i+1)*RTPPaddingMaxPayloadSize] @@ -2193,7 +2194,7 @@ func (d *DownTrack) retransmitPacket(epm *extPacketMeta, sourcePkt []byte, isPro return 0, errPayloadOverflow } hdr := RTPHeaderFactory.Get().(*rtp.Header) - *hdr = rtp.Header{ + initPooledRTPHeader(hdr, rtp.Header{ Version: pkt.Header.Version, Padding: false, Marker: epm.marker, @@ -2201,7 +2202,7 @@ func (d *DownTrack) retransmitPacket(epm *extPacketMeta, sourcePkt []byte, isPro SequenceNumber: epm.targetSeqNo, Timestamp: epm.timestamp, SSRC: d.ssrc, - } + }) rtxOffset := 0 var rtxExtSequenceNumber uint64 if rtxPT := d.payloadTypeRTX.Load(); rtxPT != 0 && d.ssrcRTX != 0 { @@ -2433,7 +2434,7 @@ func (d *DownTrack) WriteProbePackets(bytesToSend int, usePadding bool) int { for i := range num { rtxExtSequenceNumber := d.rtxSequenceNumber.Inc() hdr := RTPHeaderFactory.Get().(*rtp.Header) - *hdr = rtp.Header{ + initPooledRTPHeader(hdr, rtp.Header{ Version: 2, Padding: true, PaddingSize: byte(RTPPaddingMaxPayloadSize), @@ -2442,7 +2443,7 @@ func (d *DownTrack) WriteProbePackets(bytesToSend int, usePadding bool) int { SequenceNumber: uint16(rtxExtSequenceNumber), Timestamp: 0, SSRC: d.ssrcRTX, - } + }) d.addDummyExtensions(hdr) payload := payloads[i*RTPPaddingMaxPayloadSize : (i+1)*RTPPaddingMaxPayloadSize : (i+1)*RTPPaddingMaxPayloadSize] diff --git a/pkg/sfu/forwarder.go b/pkg/sfu/forwarder.go index 6d9069151..5dfc69038 100644 --- a/pkg/sfu/forwarder.go +++ b/pkg/sfu/forwarder.go @@ -201,12 +201,18 @@ type TranslationParams struct { rtp TranslationParamsRTP ddBytes []byte incomingHeaderSize int - codecBytes []byte + codecBytes [codecmunger.MaxHeaderSize]byte + codecBytesLen int marker bool // end of the svc spatial layer frame isEndOfLayerFrame bool } +// codecHeader returns the munged codec header to prepend to the payload. +func (tp *TranslationParams) codecHeader() []byte { + return tp.codecBytes[:tp.codecBytesLen] +} + // ------------------------------------------------------------------- type refInfo struct { @@ -2223,7 +2229,7 @@ func (f *Forwarder) getTranslationParamsVideo(extPkt *buffer.ExtPacket, layer in func (f *Forwarder) translateCodecHeader(extPkt *buffer.ExtPacket, tp *TranslationParams) error { // codec specific forwarding check and any needed packet munging tl := f.vls.SelectTemporal(extPkt) - inputSize, codecBytes, err := f.codecMunger.UpdateAndGet( + inputSize, codecBytes, codecBytesLen, err := f.codecMunger.UpdateAndGet( extPkt, tp.rtp.snOrdering == SequenceNumberOrderingOutOfOrder, tp.rtp.snOrdering == SequenceNumberOrderingGap, @@ -2243,6 +2249,7 @@ func (f *Forwarder) translateCodecHeader(extPkt *buffer.ExtPacket, tp *Translati } tp.incomingHeaderSize = inputSize tp.codecBytes = codecBytes + tp.codecBytesLen = codecBytesLen return nil } diff --git a/pkg/sfu/forwarder_test.go b/pkg/sfu/forwarder_test.go index b9d8170f8..defec62db 100644 --- a/pkg/sfu/forwarder_test.go +++ b/pkg/sfu/forwarder_test.go @@ -27,6 +27,7 @@ import ( "github.com/livekit/protocol/logger" "github.com/livekit/livekit-server/pkg/sfu/buffer" + "github.com/livekit/livekit-server/pkg/sfu/codecmunger" dd "github.com/livekit/livekit-server/pkg/sfu/rtpextension/dependencydescriptor" "github.com/livekit/livekit-server/pkg/sfu/testutils" ) @@ -1598,7 +1599,8 @@ func TestForwarderGetTranslationParamsVideo(t *testing.T) { extTimestamp: 0xabcdef, }, incomingHeaderSize: 6, - codecBytes: marshalledVP8, + codecBytes: codecHeaderOf(marshalledVP8), + codecBytesLen: len(marshalledVP8), marker: true, } actualTP, err = f.GetTranslationParams(extPkt, 0) @@ -1677,7 +1679,8 @@ func TestForwarderGetTranslationParamsVideo(t *testing.T) { extTimestamp: 0xabcdef, }, incomingHeaderSize: 6, - codecBytes: marshalledVP8, + codecBytes: codecHeaderOf(marshalledVP8), + codecBytesLen: len(marshalledVP8), } actualTP, err = f.GetTranslationParams(extPkt, 0) require.NoError(t, err) @@ -1731,7 +1734,8 @@ func TestForwarderGetTranslationParamsVideo(t *testing.T) { extTimestamp: 0xabcdef, }, incomingHeaderSize: 6, - codecBytes: marshalledVP8, + codecBytes: codecHeaderOf(marshalledVP8), + codecBytesLen: len(marshalledVP8), } actualTP, err = f.GetTranslationParams(extPkt, 0) require.NoError(t, err) @@ -1819,7 +1823,8 @@ func TestForwarderGetTranslationParamsVideo(t *testing.T) { extTimestamp: 0xabcdef, }, incomingHeaderSize: 6, - codecBytes: marshalledVP8, + codecBytes: codecHeaderOf(marshalledVP8), + codecBytesLen: len(marshalledVP8), } actualTP, err = f.GetTranslationParams(extPkt, 0) require.NoError(t, err) @@ -1918,7 +1923,8 @@ func TestForwarderGetTranslationParamsVideo(t *testing.T) { extTimestamp: 0xabcdf0, }, incomingHeaderSize: 5, - codecBytes: marshalledVP8, + codecBytes: codecHeaderOf(marshalledVP8), + codecBytesLen: len(marshalledVP8), } actualTP, err = f.GetTranslationParams(extPkt, 1) require.NoError(t, err) @@ -2249,3 +2255,9 @@ func TestForwarderIsEndOfLayerFrame(t *testing.T) { DependencyDescriptor: &buffer.ExtDependencyDescriptor{}, })) } + +func codecHeaderOf(b []byte) [codecmunger.MaxHeaderSize]byte { + var hdr [codecmunger.MaxHeaderSize]byte + copy(hdr[:], b) + return hdr +} diff --git a/pkg/sfu/pacer/base.go b/pkg/sfu/pacer/base.go index 96e3be250..b5a35e9e8 100644 --- a/pkg/sfu/pacer/base.go +++ b/pkg/sfu/pacer/base.go @@ -58,7 +58,11 @@ func (b *Base) TimeSinceLastSentPacket() time.Duration { func (b *Base) SendPacket(p *Packet) (int, error) { defer func() { if p.HeaderPool != nil && p.Header != nil { + // keep Extensions capacity so the next user does not allocate on SetExtension + exts := p.Header.Extensions + clear(exts) *p.Header = rtp.Header{} + p.Header.Extensions = exts[:0] p.HeaderPool.Put(p.Header) } @@ -95,12 +99,11 @@ func (b *Base) patchRTPHeaderExtensions(p *Packet) error { absSendTimeExt := rtp.AbsSendTimeExtension{ Timestamp: uint64(mediatransportutil.ToNtpTime(sendingAt) >> 14), } - absSendTimeBytes, err := absSendTimeExt.Marshal() - if err != nil { + if _, err := absSendTimeExt.MarshalTo(p.absSendTimeBuf[:]); err != nil { return err } - if err = p.Header.SetExtension(p.AbsSendTimeExtID, absSendTimeBytes); err != nil { + if err := p.Header.SetExtension(p.AbsSendTimeExtID, p.absSendTimeBuf[:]); err != nil { return err } @@ -119,12 +122,11 @@ func (b *Base) patchRTPHeaderExtensions(p *Packet) error { twccExt := rtp.TransportCCExtension{ TransportSequence: twccSN, } - twccExtBytes, err := twccExt.Marshal() - if err != nil { + if _, err := twccExt.MarshalTo(p.twccBuf[:]); err != nil { return err } - if err = p.Header.SetExtension(p.TransportWideExtID, twccExtBytes); err != nil { + if err := p.Header.SetExtension(p.TransportWideExtID, p.twccBuf[:]); err != nil { return err } diff --git a/pkg/sfu/pacer/base_test.go b/pkg/sfu/pacer/base_test.go new file mode 100644 index 000000000..3842203a6 --- /dev/null +++ b/pkg/sfu/pacer/base_test.go @@ -0,0 +1,87 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package pacer + +import ( + "sync" + "testing" + + "github.com/pion/rtp" + "github.com/stretchr/testify/require" + + "github.com/livekit/protocol/logger" + + "github.com/livekit/livekit-server/pkg/sfu/bwe" + "github.com/livekit/livekit-server/pkg/sfu/utils" +) + +// bwe.NullBWE is meant to be embedded and lacks Type +type nullBWE struct{ *bwe.NullBWE } + +func (nullBWE) Type() bwe.BWEType { return bwe.BWETypeNone } + +// records extension sizes without allocating so AllocsPerRun measures only SendPacket +type extCheckWriter struct { + writes int + absSendTimeLen int + transportWideLen int +} + +func (w *extCheckWriter) WriteRTP(header *rtp.Header, _ []byte) (int, error) { + w.writes++ + w.absSendTimeLen += len(header.GetExtension(1)) + w.transportWideLen += len(header.GetExtension(2)) + return 0, nil +} + +func (w *extCheckWriter) Write(_ []byte) (int, error) { return 0, nil } + +// SendPacket patches abs-send-time and transport-cc into a pooled header, +// this must not allocate per packet +func TestSendPacketHeaderExtensionsNoAlloc(t *testing.T) { + b := NewBase(logger.GetLogger(), nullBWE{&bwe.NullBWE{}}) + headerPool := &sync.Pool{New: func() any { return &rtp.Header{} }} + w := &extCheckWriter{} + payload := make([]byte, 100) + + send := func() { + hdr := headerPool.Get().(*rtp.Header) + exts := hdr.Extensions[:0] + *hdr = rtp.Header{Version: 2, SequenceNumber: 1, Timestamp: 2, SSRC: 3} + hdr.Extensions = exts + + p := PacketFactory.Get().(*Packet) + *p = Packet{ + Header: hdr, + HeaderPool: headerPool, + HeaderSize: hdr.MarshalSize(), + Payload: payload, + AbsSendTimeExtID: 1, + TransportWideExtID: 2, + WriteStream: w, + } + _, err := b.SendPacket(p) + require.NoError(t, err) + } + + send() // warm the pools; AllocsPerRun also runs once before measuring + allocs := testing.AllocsPerRun(1000, send) + require.Equal(t, 1002, w.writes) + require.Equal(t, 3*w.writes, w.absSendTimeLen) + require.Equal(t, 2*w.writes, w.transportWideLen) + if !utils.RaceEnabled { + require.Equal(t, 0.0, allocs, "allocations per SendPacket") + } +} diff --git a/pkg/sfu/pacer/pacer.go b/pkg/sfu/pacer/pacer.go index f23212519..199a8bf17 100644 --- a/pkg/sfu/pacer/pacer.go +++ b/pkg/sfu/pacer/pacer.go @@ -54,6 +54,11 @@ type Packet struct { WriteStream webrtc.TrackLocalWriter Pool *sync.Pool PoolEntity *[]byte + + // per-packet scratch for header extensions patched at send time, + // the Packet is owned by one send until SendPacket returns + absSendTimeBuf [3]byte + twccBuf [2]byte } type Pacer interface { diff --git a/pkg/sfu/receiver_base.go b/pkg/sfu/receiver_base.go index 80874a893..617d7fc64 100644 --- a/pkg/sfu/receiver_base.go +++ b/pkg/sfu/receiver_base.go @@ -962,12 +962,9 @@ func (r *ReceiverBase) forwardRTP( continue } - var writeCount atomic.Int32 - r.downTrackSpreader.Broadcast(func(dt TrackSender) { - writeCount.Add(dt.WriteRTP(extPkt, spatialLayer)) - }) + writeCount := sfuutils.BroadcastRTP(r.downTrackSpreader, extPkt, spatialLayer) if rt := r.loadREDTransformer(); rt != nil { - writeCount.Add(rt.ForwardRTP(extPkt, spatialLayer)) + writeCount += rt.ForwardRTP(extPkt, spatialLayer) } // track delay/jitter @@ -977,13 +974,13 @@ func (r *ReceiverBase) forwardRTP( // delivered back-to-back) which the single forwarder goroutine drains // serially, inflating the measured transit for the tail of the burst. That // reflects loss recovery rather than steady-state forwarding health. - if writeCount.Load() > 0 && r.forwardStats != nil && !extPkt.IsBuffered && !extPkt.IsOutOfOrder { + if writeCount > 0 && r.forwardStats != nil && !extPkt.IsBuffered && !extPkt.IsOutOfOrder { if latency, isHigh := r.forwardStats.Update(extPkt.Arrival, mono.UnixNano()); isHigh { r.params.Logger.Debugw( "high forwarding latency", "latency", time.Duration(latency), "queuingLatency", time.Duration(dequeuedAt-extPkt.Arrival), - "writeCount", writeCount.Load(), + "writeCount", writeCount, "isOutOfOrder", extPkt.IsOutOfOrder, "layer", layer, ) diff --git a/pkg/sfu/redprimaryreceiver.go b/pkg/sfu/redprimaryreceiver.go index 17b639abb..a34ef8586 100644 --- a/pkg/sfu/redprimaryreceiver.go +++ b/pkg/sfu/redprimaryreceiver.go @@ -50,6 +50,10 @@ type RedPrimaryReceiver struct { // bitset for upstream packet receive history [lastSeq-8, lastSeq-1], bit 1 represents packet received pktHistory byte + + // forwarded packet, reused since ForwardRTP runs on one goroutine and + // down tracks do not keep the packet past WriteRTP + sendExtPkt buffer.ExtPacket } func NewRedPrimaryReceiver(receiver TrackReceiver, dsp utils.DownTrackSpreaderParams) REDTransformer { @@ -73,11 +77,7 @@ func (r *RedPrimaryReceiver) ForwardRTP(pkt *buffer.ExtPacket, spatialLayer int3 if pkt.Packet.PayloadType != r.redPT { // forward non-red packet directly - var writeCount atomic.Int32 - r.downTrackSpreader.Broadcast(func(dt TrackSender) { - writeCount.Add(dt.WriteRTP(pkt, spatialLayer)) - }) - return writeCount.Load() + return utils.BroadcastRTP(r.downTrackSpreader, pkt, spatialLayer) } pkts, err := r.getSendPktsFromRed(pkt.Packet) @@ -86,9 +86,10 @@ func (r *RedPrimaryReceiver) ForwardRTP(pkt *buffer.ExtPacket, spatialLayer int3 return 0 } - var writeCount atomic.Int32 + var writeCount int32 for i, sendPkt := range pkts { - pPkt := *pkt + pPkt := &r.sendExtPkt + *pPkt = *pkt if i != len(pkts)-1 { // patch extended sequence number and time stamp for all but the last packet, // last packet is the primary payload @@ -129,11 +130,9 @@ func (r *RedPrimaryReceiver) ForwardRTP(pkt *buffer.ExtPacket, spatialLayer int3 // not modify the ExtPacket.RawPacket here for performance since it is not used by the DownTrack, // otherwise it should be set to the correct value (marshal the primary rtp packet) - r.downTrackSpreader.Broadcast(func(dt TrackSender) { - writeCount.Add(dt.WriteRTP(&pPkt, spatialLayer)) - }) + writeCount += utils.BroadcastRTP(r.downTrackSpreader, pPkt, spatialLayer) } - return writeCount.Load() + return writeCount } func (r *RedPrimaryReceiver) ForwardRTCPSenderReport( diff --git a/pkg/sfu/redreceiver.go b/pkg/sfu/redreceiver.go index 959a54e96..4977878bf 100644 --- a/pkg/sfu/redreceiver.go +++ b/pkg/sfu/redreceiver.go @@ -51,6 +51,10 @@ type RedReceiver struct { closed atomic.Bool pktBuff [maxRedCount]*rtp.Packet redPayloadBuf [mtuSize]byte + // forwarded packet, reused like redPayloadBuf since ForwardRTP runs on one goroutine + // and down tracks do not keep the packet past WriteRTP + redExtPkt buffer.ExtPacket + redRtpPkt rtp.Packet } func NewRedReceiver(receiver TrackReceiver, dsp utils.DownTrackSpreaderParams) REDTransformer { @@ -73,11 +77,7 @@ func (r *RedReceiver) ForwardRTP(pkt *buffer.ExtPacket, spatialLayer int32) int3 // fallback to primary codec if payload size exceeds redundant block length if len(pkt.Packet.Payload) >= maxRedPayload { - var writeCount atomic.Int32 - r.downTrackSpreader.Broadcast(func(dt TrackSender) { - writeCount.Add(dt.WriteRTP(pkt, spatialLayer)) - }) - return writeCount.Load() + return utils.BroadcastRTP(r.downTrackSpreader, pkt, spatialLayer) } redLen, err := r.encodeRedForPrimary(pkt.Packet, r.redPayloadBuf[:]) @@ -86,19 +86,17 @@ func (r *RedReceiver) ForwardRTP(pkt *buffer.ExtPacket, spatialLayer int32) int3 return 0 } - pPkt := *pkt - redRtpPacket := *pkt.Packet + redRtpPacket := &r.redRtpPkt + *redRtpPacket = *pkt.Packet redRtpPacket.PayloadType = opusRedPT redRtpPacket.Payload = r.redPayloadBuf[:redLen] - pPkt.Packet = &redRtpPacket + pPkt := &r.redExtPkt + *pPkt = *pkt + pPkt.Packet = redRtpPacket // not modify the ExtPacket.RawPacket here for performance since it is not used by the DownTrack, // otherwise it should be set to the correct value (marshal the primary rtp packet) - var writeCount atomic.Int32 - r.downTrackSpreader.Broadcast(func(dt TrackSender) { - writeCount.Add(dt.WriteRTP(&pPkt, spatialLayer)) - }) - return writeCount.Load() + return utils.BroadcastRTP(r.downTrackSpreader, pPkt, spatialLayer) } func (r *RedReceiver) ForwardRTCPSenderReport( diff --git a/pkg/sfu/sfu.go b/pkg/sfu/sfu.go index c8760d069..1cd41d13f 100644 --- a/pkg/sfu/sfu.go +++ b/pkg/sfu/sfu.go @@ -34,3 +34,10 @@ var ( }, } ) + +// initPooledRTPHeader sets hdr to h while keeping the pooled header's +// Extensions capacity, so SetExtension does not allocate per packet. +func initPooledRTPHeader(hdr *rtp.Header, h rtp.Header) { + h.Extensions = hdr.Extensions[:0] + *hdr = h +} diff --git a/pkg/sfu/utils/downtrackspreader.go b/pkg/sfu/utils/downtrackspreader.go index e06f6c945..29ba07f2d 100644 --- a/pkg/sfu/utils/downtrackspreader.go +++ b/pkg/sfu/utils/downtrackspreader.go @@ -15,16 +15,18 @@ package utils import ( + "runtime" "sync" + "sync/atomic" "github.com/livekit/protocol/livekit" "github.com/livekit/protocol/logger" "github.com/livekit/protocol/utils" ) -type sender interface { - SubscriberID() livekit.ParticipantID -} +// 100µs is enough to amortize the overhead and provide sufficient load balancing. +// WriteRTP takes about 50µs on average, so we write to 2 down tracks per loop. +const broadcastStep = 2 type DownTrackSpreaderParams struct { Threshold int @@ -90,23 +92,26 @@ func (d *DownTrackSpreader[T]) HasDownTrack(subscriberID livekit.ParticipantID) return ok } -func (d *DownTrackSpreader[T]) Broadcast(writer func(T)) { +// snapshot returns the current down tracks and the parallelization threshold +func (d *DownTrackSpreader[T]) snapshot() ([]T, int) { d.downTrackMu.RLock() downTracks := d.downTracksShadow - threshold := uint64(d.params.Threshold) + threshold := d.params.Threshold d.downTrackMu.RUnlock() - if len(downTracks) == 0 { - return - } if threshold == 0 { threshold = 1000000 } + return downTracks, threshold +} - // 100µs is enough to amortize the overhead and provide sufficient load balancing. - // WriteRTP takes about 50µs on average, so we write to 2 down tracks per loop. - step := uint64(2) - utils.ParallelExec(downTracks, threshold, step, writer) +func (d *DownTrackSpreader[T]) Broadcast(writer func(T)) { + downTracks, threshold := d.snapshot() + if len(downTracks) == 0 { + return + } + + utils.ParallelExec(downTracks, uint64(threshold), broadcastStep, writer) } func (d *DownTrackSpreader[T]) DownTrackCount() int { @@ -127,3 +132,69 @@ func (d *DownTrackSpreader[T]) SetThreshold(threshold int) { d.params.Threshold = threshold d.downTrackMu.Unlock() } + +// ------------------------------------------------ + +type sender interface { + SubscriberID() livekit.ParticipantID +} + +type rtpWriter[P any] interface { + WriteRTP(pkt P, layer int32) int32 +} + +// rtpBroadcast is the shared state of one parallel BroadcastRTP, it carries the +// packet and layer so that no closure has to be allocated per packet +type rtpBroadcast[T rtpWriter[P], P any] struct { + downTracks []T + pkt P + layer int32 + next atomic.Uint64 + written atomic.Int32 + wg sync.WaitGroup +} + +func (b *rtpBroadcast[T, P]) run() { + defer b.wg.Done() + + var written int32 + end := uint64(len(b.downTracks)) + for { + n := b.next.Add(broadcastStep) + if n >= end+broadcastStep { + break + } + for i := n - broadcastStep; i < n && i < end; i++ { + written += b.downTracks[i].WriteRTP(b.pkt, b.layer) + } + } + b.written.Add(written) +} + +// BroadcastRTP writes pkt to every down track and returns how many accepted it. +// Below the threshold it runs on the caller without allocating; above it, the +// only allocations are the shared state and the worker funcval. +func BroadcastRTP[T interface { + sender + rtpWriter[P] +}, P any](d *DownTrackSpreader[T], pkt P, layer int32) int32 { + downTracks, threshold := d.snapshot() + if len(downTracks) < threshold { + var written int32 + for _, dt := range downTracks { + written += dt.WriteRTP(pkt, layer) + } + return written + } + + numWorkers := min(runtime.NumCPU(), len(downTracks)) + b := &rtpBroadcast[T, P]{downTracks: downTracks, pkt: pkt, layer: layer} + b.wg.Add(numWorkers) + worker := b.run + for i := 0; i < numWorkers; i++ { + go worker() + } + b.wg.Wait() + + return b.written.Load() +} diff --git a/pkg/sfu/utils/downtrackspreader_test.go b/pkg/sfu/utils/downtrackspreader_test.go new file mode 100644 index 000000000..1de2d96ce --- /dev/null +++ b/pkg/sfu/utils/downtrackspreader_test.go @@ -0,0 +1,84 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package utils + +import ( + "fmt" + "sync/atomic" + "testing" + + "github.com/stretchr/testify/require" + + "github.com/livekit/protocol/livekit" +) + +type testPacket struct{ seq int } + +type testSender struct { + id livekit.ParticipantID + writes atomic.Int32 + accept int32 +} + +func (s *testSender) SubscriberID() livekit.ParticipantID { return s.id } + +func (s *testSender) WriteRTP(pkt *testPacket, layer int32) int32 { + s.writes.Add(1) + return s.accept +} + +func newTestSpreader(numDownTracks, threshold int) (*DownTrackSpreader[*testSender], []*testSender) { + d := NewDownTrackSpreader[*testSender](DownTrackSpreaderParams{Threshold: threshold}) + senders := make([]*testSender, numDownTracks) + for i := range senders { + senders[i] = &testSender{id: livekit.ParticipantID(fmt.Sprintf("p%d", i)), accept: 1} + d.Store(senders[i]) + } + return d, senders +} + +func TestBroadcastRTP(t *testing.T) { + for _, tc := range []struct { + name string + numDownTracks int + threshold int + maxAllocs float64 + }{ + {"serial", 5, 20, 0}, + {"parallel", 50, 20, 2}, // shared state and worker funcval + } { + t.Run(tc.name, func(t *testing.T) { + d, senders := newTestSpreader(tc.numDownTracks, tc.threshold) + senders[0].accept = 0 // one down track that drops + pkt := &testPacket{} + + written := BroadcastRTP(d, pkt, 2) + require.EqualValues(t, tc.numDownTracks-1, written) + for _, s := range senders { + require.EqualValues(t, 1, s.writes.Load()) + } + + if !RaceEnabled { + allocs := testing.AllocsPerRun(1000, func() { BroadcastRTP(d, pkt, 2) }) + require.LessOrEqual(t, allocs, tc.maxAllocs) + } + }) + } +} + +func TestBroadcastRTPEmpty(t *testing.T) { + d := NewDownTrackSpreader[*testSender](DownTrackSpreaderParams{Threshold: 20}) + require.EqualValues(t, 0, BroadcastRTP(d, &testPacket{}, 0)) +} diff --git a/pkg/sfu/utils/norace.go b/pkg/sfu/utils/norace.go new file mode 100644 index 000000000..ba683f0f7 --- /dev/null +++ b/pkg/sfu/utils/norace.go @@ -0,0 +1,6 @@ +//go:build !race + +package utils + +// RaceEnabled reports whether the binary was built with the race detector. +const RaceEnabled = false diff --git a/pkg/sfu/utils/race.go b/pkg/sfu/utils/race.go new file mode 100644 index 000000000..72223acef --- /dev/null +++ b/pkg/sfu/utils/race.go @@ -0,0 +1,8 @@ +//go:build race + +package utils + +// RaceEnabled reports whether the binary was built with the race detector. +// sync.Pool drops a quarter of returned items under it, so allocation counts +// on pooled paths are not meaningful. +const RaceEnabled = true diff --git a/version/version.go b/version/version.go index 29d306b8c..0486d6293 100644 --- a/version/version.go +++ b/version/version.go @@ -14,4 +14,4 @@ package version -const Version = "1.13.6" +const Version = "1.13.7"