From 8d11efdfcd4220092b6ac7b8a21af28526da5a6b Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Mon, 14 Sep 2026 10:20:40 +0530 Subject: [PATCH 1/8] Release v1.13.7. (#4866) Co-authored-by: Claude Opus 5 (1M context) --- CHANGELOG.md | 42 ++++++++++++++++++++++++++++++++++++++++++ version/version.go | 2 +- 2 files changed, 43 insertions(+), 1 deletion(-) 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/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" From 776185e4a5b74210f048327f57880311ce3033e0 Mon Sep 17 00:00:00 2001 From: "renovate[bot]" <29139614+renovate[bot]@users.noreply.github.com> Date: Sun, 13 Sep 2026 22:38:49 -0700 Subject: [PATCH 2/8] Update module github.com/gotesttools/gotestfmt/v2 to v2.5.0 (#4865) Generated by renovateBot Co-authored-by: renovate[bot] <29139614+renovate[bot]@users.noreply.github.com> --- go.mod | 2 +- go.sum | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/go.mod b/go.mod index b2bdab9dd..e131ed8c0 100644 --- a/go.mod +++ b/go.mod @@ -74,7 +74,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 diff --git a/go.sum b/go.sum index 49d89faf0..08d64c4d9 100644 --- a/go.sum +++ b/go.sum @@ -105,8 +105,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= From 08c5433449c106943231037d0b892895f6e92d3b Mon Sep 17 00:00:00 2001 From: "renovate[bot]" <29139614+renovate[bot]@users.noreply.github.com> Date: Sun, 13 Sep 2026 23:04:12 -0700 Subject: [PATCH 3/8] Update module github.com/urfave/cli/v3 to v3.11.0 (#4867) Generated by renovateBot Co-authored-by: renovate[bot] <29139614+renovate[bot]@users.noreply.github.com> --- go.mod | 2 +- go.sum | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/go.mod b/go.mod index e131ed8c0..cc7cdadf6 100644 --- a/go.mod +++ b/go.mod @@ -143,7 +143,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 08d64c4d9..42e1d8aca 100644 --- a/go.sum +++ b/go.sum @@ -307,8 +307,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= From 0392ca90b40503438c0022cfc895955e05d65457 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Tue, 15 Sep 2026 21:34:10 +0530 Subject: [PATCH 4/8] sfu: return VP8 munged header by value to avoid per-packet allocation (#4871) CodecMunger.UpdateAndGet returned a freshly allocated slice for every forwarded packet on every down track. Return the header in a fixed array with a length instead and marshal into it with MarshalTo. TranslationParams carries the array; WriteRTP slices it locally. UpdateAndGet now allocates zero times per call (was one). Escape analysis confirms TranslationParams stays on the stack. Co-authored-by: Claude Fable 5.1 --- pkg/sfu/codecmunger/codecmunger.go | 7 ++++++- pkg/sfu/codecmunger/null.go | 4 ++-- pkg/sfu/codecmunger/vp8.go | 21 ++++++++++---------- pkg/sfu/codecmunger/vp8_test.go | 32 +++++++++++++++--------------- pkg/sfu/downtrack.go | 9 +++++---- pkg/sfu/forwarder.go | 11 ++++++++-- pkg/sfu/forwarder_test.go | 22 +++++++++++++++----- 7 files changed, 66 insertions(+), 40 deletions(-) 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..e48d717d7 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 { @@ -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, 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 +} From 61a31d31ff80ba29c8b8e7530840b87736f49ab9 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Tue, 15 Sep 2026 21:48:24 +0530 Subject: [PATCH 5/8] sfu: patch pacer header extensions without per-packet allocation (#4872) SendPacket allocated four times per down track write: a 3 byte slice for abs-send-time, a 2 byte slice for transport-cc, and two appends to a nil Extensions slice because the pooled header was reset to a zero value. Give pacer.Packet fixed scratch arrays for the two extensions and marshal into them. The Packet is pooled and owned by one send until SendPacket returns, so nothing is shared across down tracks. Keep the pooled header's Extensions capacity across reuse, both when returning it to the pool and when a down track initializes it. Adds a test that asserts zero allocations per SendPacket. Co-authored-by: Claude Fable 5.1 --- pkg/sfu/downtrack.go | 16 ++++---- pkg/sfu/pacer/base.go | 14 ++++--- pkg/sfu/pacer/base_test.go | 84 ++++++++++++++++++++++++++++++++++++++ pkg/sfu/pacer/pacer.go | 5 +++ pkg/sfu/sfu.go | 7 ++++ 5 files changed, 112 insertions(+), 14 deletions(-) create mode 100644 pkg/sfu/pacer/base_test.go diff --git a/pkg/sfu/downtrack.go b/pkg/sfu/downtrack.go index e48d717d7..d83c020f4 100644 --- a/pkg/sfu/downtrack.go +++ b/pkg/sfu/downtrack.go @@ -1069,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, @@ -1077,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 { @@ -1281,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), @@ -1290,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] @@ -2194,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, @@ -2202,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 { @@ -2434,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), @@ -2443,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/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..3b68fa168 --- /dev/null +++ b/pkg/sfu/pacer/base_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 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" +) + +// 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) + 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/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 +} From 482fbe2cf7d030e4b79d2995f0e94ff6a9670100 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Wed, 16 Sep 2026 02:59:09 +0530 Subject: [PATCH 6/8] sfu: skip pacer allocation assertion under the race detector (#4874) sync.Pool drops a quarter of returned items when built with -race, so the pooled send path averages about one allocation per packet in CI and the zero-allocation assertion flips between passing and failing. Keep the extension checks and skip only the allocation count under race. Co-authored-by: Claude Fable 5.1 --- pkg/sfu/pacer/base_test.go | 5 ++++- pkg/sfu/utils/norace.go | 6 ++++++ pkg/sfu/utils/race.go | 8 ++++++++ 3 files changed, 18 insertions(+), 1 deletion(-) create mode 100644 pkg/sfu/utils/norace.go create mode 100644 pkg/sfu/utils/race.go diff --git a/pkg/sfu/pacer/base_test.go b/pkg/sfu/pacer/base_test.go index 3b68fa168..3842203a6 100644 --- a/pkg/sfu/pacer/base_test.go +++ b/pkg/sfu/pacer/base_test.go @@ -24,6 +24,7 @@ import ( "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 @@ -80,5 +81,7 @@ func TestSendPacketHeaderExtensionsNoAlloc(t *testing.T) { require.Equal(t, 1002, w.writes) require.Equal(t, 3*w.writes, w.absSendTimeLen) require.Equal(t, 2*w.writes, w.transportWideLen) - require.Equal(t, 0.0, allocs, "allocations per SendPacket") + if !utils.RaceEnabled { + require.Equal(t, 0.0, allocs, "allocations per SendPacket") + } } 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 From 4ea0facf856aa19ea601a896fc3db941420fe151 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Wed, 16 Sep 2026 03:09:40 +0530 Subject: [PATCH 7/8] Update pion/sdp to get lesser SDP retention (#4873) --- go.mod | 2 +- go.sum | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/go.mod b/go.mod index cc7cdadf6..1e0088fe5 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 diff --git a/go.sum b/go.sum index 42e1d8aca..df7274bd3 100644 --- a/go.sum +++ b/go.sum @@ -253,8 +253,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= From 48700da3f3f2c3b514f673308bed861ec09fd729 Mon Sep 17 00:00:00 2001 From: Raja Subramanian Date: Wed, 16 Sep 2026 11:53:54 +0530 Subject: [PATCH 8/8] sfu: broadcast RTP to down tracks without per-packet allocations (#4875) * sfu: broadcast RTP to down tracks without a per-packet closure Every forwarded packet allocated a closure capturing the packet and layer, and an atomic write counter that escaped with it, even for tracks with a single subscriber. Add BroadcastRTP, which carries the packet and layer in the worker state and sums the writes per worker. The serial path no longer allocates; the parallel path allocates the shared state and the worker funcval only. Co-Authored-By: Claude Fable 5.1 * sfu: reuse the forwarded packet copy in the RED receivers Both RED receivers copied the ExtPacket (and the opus one the rtp.Packet) per forwarded packet, and the copies moved to the heap because the broadcast takes their address. Keep them on the receiver like the existing redPayloadBuf: ForwardRTP runs on one goroutine and down tracks do not keep the packet past WriteRTP. Co-Authored-By: Claude Fable 5.1 * sfu: skip BroadcastRTP allocation assertion under the race detector Co-Authored-By: Claude Fable 5.1 * sfu: group BroadcastRTP and its interfaces after the spreader methods Co-Authored-By: Claude Fable 5.1 --------- Co-authored-by: Claude Fable 5.1 --- pkg/sfu/receiver_base.go | 11 ++- pkg/sfu/redprimaryreceiver.go | 21 +++--- pkg/sfu/redreceiver.go | 24 +++---- pkg/sfu/utils/downtrackspreader.go | 95 +++++++++++++++++++++---- pkg/sfu/utils/downtrackspreader_test.go | 84 ++++++++++++++++++++++ 5 files changed, 192 insertions(+), 43 deletions(-) create mode 100644 pkg/sfu/utils/downtrackspreader_test.go 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/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)) +}