mirror of
https://github.com/livekit/livekit.git
synced 2026-09-17 01:34:51 +00:00
Merge remote-tracking branch 'origin/master' into pr-4779-local
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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=
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
|
||||
+11
-10
@@ -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) {
|
||||
|
||||
@@ -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)
|
||||
|
||||
+13
-12
@@ -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]
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
|
||||
@@ -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(
|
||||
|
||||
+11
-13
@@ -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(
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
|
||||
@@ -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))
|
||||
}
|
||||
@@ -0,0 +1,6 @@
|
||||
//go:build !race
|
||||
|
||||
package utils
|
||||
|
||||
// RaceEnabled reports whether the binary was built with the race detector.
|
||||
const RaceEnabled = false
|
||||
@@ -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
|
||||
+1
-1
@@ -14,4 +14,4 @@
|
||||
|
||||
package version
|
||||
|
||||
const Version = "1.13.6"
|
||||
const Version = "1.13.7"
|
||||
|
||||
Reference in New Issue
Block a user