diff --git a/pkg/config/config.go b/pkg/config/config.go index 3ca89107e..04a4a0f6f 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -186,6 +186,18 @@ type RTCConfig struct { // enable rtp stream restart detection for published tracks EnableRTPStreamRestartDetection bool `yaml:"enable_rtp_stream_restart_detection,omitempty"` + + // strip codec/extension/ssrc attributes from inactive media sections in client offers + // before they reach pion, bounds memory retained for clients that add a media section per publish + ShrinkInactiveMediaSections ShrinkInactiveMediaSectionsConfig `yaml:"shrink_inactive_media_sections,omitempty"` +} + +type ShrinkInactiveMediaSectionsConfig struct { + Enabled bool `yaml:"enabled,omitempty"` + // offers smaller than this are left untouched + MinSDPSize int `yaml:"min_sdp_size,omitempty"` + // video sections are always eligible, audio only when set + IncludeAudio bool `yaml:"include_audio,omitempty"` } type TURNServer struct { @@ -539,6 +551,7 @@ var DefaultConfig = Config{ PacketBufferSize: 500, PacketBufferSizeVideo: 500, PacketBufferSizeAudio: 200, + ShrinkInactiveMediaSections: ShrinkInactiveMediaSectionsConfig{MinSDPSize: 20_000}, PLIThrottle: sfu.DefaultPLIThrottleConfig, DatachannelDataTrackTargetLatency: 100 * time.Millisecond, CongestionControl: CongestionControlConfig{ diff --git a/pkg/rtc/participant.go b/pkg/rtc/participant.go index 01b14479d..7ecfd2042 100644 --- a/pkg/rtc/participant.go +++ b/pkg/rtc/participant.go @@ -240,6 +240,7 @@ type ParticipantParams struct { EnableStartAtDesiredQuality bool MigrationWaitDuration time.Duration ExcludeIPv6LocalCandidates bool + ShrinkInactiveMediaSections config.ShrinkInactiveMediaSectionsConfig EnableWarp bool } @@ -2129,6 +2130,7 @@ func (p *ParticipantImpl) setupTransportManager() error { FireOnTrackBySdp: p.params.FireOnTrackBySdp, EnableDataTracks: p.params.EnableDataTracks, ExcludeIPv6LocalCandidates: p.params.ExcludeIPv6LocalCandidates, + ShrinkInactiveMediaSections: p.params.ShrinkInactiveMediaSections, EnableWarp: p.params.EnableWarp, } if p.params.SyncStreams && p.params.PlayoutDelay.GetEnabled() && p.params.ClientInfo.isFirefox() { diff --git a/pkg/rtc/transport.go b/pkg/rtc/transport.go index 2f9f7d440..dbbfcc97a 100644 --- a/pkg/rtc/transport.go +++ b/pkg/rtc/transport.go @@ -320,6 +320,7 @@ type TransportParams struct { AllowPlayoutDelay bool UseOneShotSignallingMode bool ExcludeIPv6LocalCandidates bool + ShrinkInactiveMediaSections config.ShrinkInactiveMediaSectionsConfig FireOnTrackBySdp bool DataChannelMaxBufferedAmount uint64 DatachannelSlowThreshold int @@ -2534,6 +2535,23 @@ func (t *PCTransport) filterCandidates(sd webrtc.SessionDescription, preferTCP, m.Attributes = filterAttributes(m.Attributes) } + shrinkConfig := t.params.ShrinkInactiveMediaSections + if !isLocal && sd.Type == webrtc.SDPTypeOffer && shrinkConfig.Enabled && len(sd.SDP) >= shrinkConfig.MinSDPSize { + knownMids := make(map[string]bool) + for _, tr := range t.pc.GetTransceivers() { + if mid := tr.Mid(); mid != "" { + knownMids[mid] = true + } + } + if shrunk := shrinkInactiveMediaSections(parsed, knownMids, shrinkConfig.IncludeAudio); shrunk > 0 { + t.params.Logger.Debugw( + "shrunk inactive media sections in remote offer", + "count", shrunk, + "sdpSize", len(sd.SDP), + ) + } + } + bytes, err := parsed.Marshal() if err != nil { t.params.Logger.Warnw("could not marshal SDP to filter candidates", err) @@ -2543,6 +2561,75 @@ func (t *PCTransport) filterCandidates(sd webrtc.SessionDescription, preferTCP, return sd } +// attributes pion still needs from an inactive media section it already has a transceiver for +var inactiveMediaSectionAttributes = map[string]bool{ + sdp.AttrKeyMID: true, + sdp.AttrKeyInactive: true, + sdp.AttrKeyRTCPMux: true, + "ice-ufrag": true, + "ice-pwd": true, + "ice-options": true, + "fingerprint": true, + "setup": true, + "bundle-only": true, +} + +// shrinkInactiveMediaSections drops codec, header extension, msid and ssrc +// attributes from inactive video (and optionally audio) sections whose mid +// already maps to a transceiver, keeping one payload type with its rtpmap so +// pion's codec lookup still succeeds. pion keeps substrings of every remote offer it parses for as +// long as the transceivers created from it live, so a client that adds a media +// section per publish and leaves the old ones inactive makes the retained text +// grow with the square of its publish count. Shrinking the sections pion no +// longer needs keeps the retained text small. Returns the number of sections +// shrunk. +func shrinkInactiveMediaSections(parsed *sdp.SessionDescription, knownMids map[string]bool, includeAudio bool) int { + shrunk := 0 + for _, m := range parsed.MediaDescriptions { + if m.MediaName.Media != "video" && !(includeAudio && m.MediaName.Media == "audio") { + continue + } + if _, inactive := m.Attribute(sdp.AttrKeyInactive); !inactive { + continue + } + if mid, ok := m.Attribute(sdp.AttrKeyMID); !ok || !knownMids[mid] { + continue + } + + // first payload type that has an rtpmap + var rtpmap *sdp.Attribute + keepFormat := "" + for _, f := range m.MediaName.Formats { + for i := range m.Attributes { + a := &m.Attributes[i] + if a.Key == "rtpmap" && strings.HasPrefix(a.Value, f+" ") { + rtpmap = a + keepFormat = f + break + } + } + if rtpmap != nil { + break + } + } + if rtpmap == nil { + continue + } + + kept := make([]sdp.Attribute, 0, len(inactiveMediaSectionAttributes)+1) + for _, a := range m.Attributes { + if inactiveMediaSectionAttributes[a.Key] { + kept = append(kept, a) + } + } + kept = append(kept, *rtpmap) + m.Attributes = kept + m.MediaName.Formats = []string{keepFormat} + shrunk++ + } + return shrunk +} + func (t *PCTransport) clearSignalStateCheckTimer() { if t.signalStateCheckTimer != nil { t.signalStateCheckTimer.Stop() diff --git a/pkg/rtc/transport_test.go b/pkg/rtc/transport_test.go index 3c0ea425e..e3b29a388 100644 --- a/pkg/rtc/transport_test.go +++ b/pkg/rtc/transport_test.go @@ -26,6 +26,7 @@ import ( "github.com/stretchr/testify/require" "go.uber.org/atomic" + "github.com/livekit/livekit-server/pkg/config" "github.com/livekit/livekit-server/pkg/rtc/transport" "github.com/livekit/livekit-server/pkg/rtc/transport/transportfakes" "github.com/livekit/livekit-server/pkg/testutils" @@ -803,3 +804,163 @@ func TestSinglePCAnswerStripsSubscribeOnlyCodecsFromRecvSide(t *testing.T) { "answer must not advertise H.264 in recv-side m-section: %s", a.Value) } } + +// --------------------------------------------------------------------------- +// inactive media section shrinking + +type testOfferSection struct { + mid int + inactive bool + audio bool +} + +// testOffer builds a Chrome-like client offer with one video section per entry, +// VP8+rtx, header extensions and ssrc lines, bundled with shared ICE credentials. +func testOffer(sections []testOfferSection) string { + var b strings.Builder + b.WriteString("v=0\r\no=- 4611731400430051336 2 IN IP4 127.0.0.1\r\ns=-\r\nt=0 0\r\na=group:BUNDLE") + for _, s := range sections { + fmt.Fprintf(&b, " %d", s.mid) + } + b.WriteString("\r\na=extmap-allow-mixed\r\na=msid-semantic: WMS\r\n") + for _, s := range sections { + if s.audio { + b.WriteString("m=audio 9 UDP/TLS/RTP/SAVPF 111\r\nc=IN IP4 0.0.0.0\r\na=rtcp:9 IN IP4 0.0.0.0\r\n") + } else { + b.WriteString("m=video 9 UDP/TLS/RTP/SAVPF 96 97\r\nc=IN IP4 0.0.0.0\r\na=rtcp:9 IN IP4 0.0.0.0\r\n") + } + b.WriteString("a=ice-ufrag:AbCdEf01\r\na=ice-pwd:0123456789abcdefghijklmnop\r\na=ice-options:trickle\r\n") + b.WriteString("a=fingerprint:sha-256 6B:8B:F0:65:5F:78:E2:51:3B:AC:6F:F3:3F:46:1B:35:DC:B8:5F:64:1A:24:C2:43:F0:A1:58:D0:A1:2C:19:08\r\na=setup:actpass\r\n") + fmt.Fprintf(&b, "a=mid:%d\r\n", s.mid) + b.WriteString("a=extmap:2 http://www.webrtc.org/experiments/rtp-hdrext/abs-send-time\r\n") + b.WriteString("a=extmap:4 http://www.ietf.org/id/draft-holmer-rmcat-transport-wide-cc-extensions-01\r\n") + if s.inactive { + b.WriteString("a=inactive\r\n") + } else { + b.WriteString("a=sendonly\r\n") + } + fmt.Fprintf(&b, "a=msid:stream%d track%d\r\na=rtcp-mux\r\na=rtcp-rsize\r\n", s.mid, s.mid) + if s.audio { + b.WriteString("a=rtpmap:111 opus/48000/2\r\na=rtcp-fb:111 transport-cc\r\na=fmtp:111 minptime=10;useinbandfec=1\r\n") + } else { + b.WriteString("a=rtpmap:96 VP8/90000\r\na=rtcp-fb:96 goog-remb\r\na=rtcp-fb:96 transport-cc\r\na=rtcp-fb:96 nack\r\na=rtcp-fb:96 nack pli\r\n") + b.WriteString("a=rtpmap:97 rtx/90000\r\na=fmtp:97 apt=96\r\n") + } + s1, s2 := 1000000+2*s.mid, 1000001+2*s.mid + fmt.Fprintf(&b, "a=ssrc-group:FID %d %d\r\n", s1, s2) + fmt.Fprintf(&b, "a=ssrc:%d cname:c%d\r\na=ssrc:%d msid:stream%d track%d\r\n", s1, s.mid, s1, s.mid, s.mid) + fmt.Fprintf(&b, "a=ssrc:%d cname:c%d\r\na=ssrc:%d msid:stream%d track%d\r\n", s2, s.mid, s2, s.mid, s.mid) + } + return b.String() +} + +func TestShrinkInactiveMediaSections(t *testing.T) { + offer := webrtc.SessionDescription{ + Type: webrtc.SDPTypeOffer, + SDP: testOffer([]testOfferSection{ + {mid: 0, inactive: true}, // known mid, shrunk + {mid: 1, inactive: true}, // unknown mid, left alone so pion can create the transceiver + {mid: 2, inactive: false}, // active, left alone + {mid: 3, inactive: true, audio: true}, // known mid, audio, only shrunk when included + }), + } + parsed, err := offer.Unmarshal() + require.NoError(t, err) + before := make([]int, len(parsed.MediaDescriptions)) + for i, m := range parsed.MediaDescriptions { + before[i] = len(m.Attributes) + } + knownMids := map[string]bool{"0": true, "2": true, "3": true} + + require.Equal(t, 1, shrinkInactiveMediaSections(parsed, knownMids, false)) + + shrunk := parsed.MediaDescriptions[0] + require.Equal(t, []string{"96"}, shrunk.MediaName.Formats) + for _, a := range shrunk.Attributes { + require.True(t, inactiveMediaSectionAttributes[a.Key] || a.Key == "rtpmap", "unexpected attribute %s", a.Key) + } + rtpmap, ok := shrunk.Attribute("rtpmap") + require.True(t, ok) + require.Equal(t, "96 VP8/90000", rtpmap) + for _, key := range []string{sdp.AttrKeyMID, sdp.AttrKeyInactive, "ice-ufrag", "ice-pwd", "fingerprint", "setup", sdp.AttrKeyRTCPMux} { + _, ok := shrunk.Attribute(key) + require.True(t, ok, "missing %s", key) + } + require.Less(t, len(shrunk.Attributes), before[0]) + + require.Equal(t, before[1], len(parsed.MediaDescriptions[1].Attributes)) + require.Equal(t, before[2], len(parsed.MediaDescriptions[2].Attributes)) + require.Equal(t, before[3], len(parsed.MediaDescriptions[3].Attributes)) + require.Equal(t, []string{"96", "97"}, parsed.MediaDescriptions[1].MediaName.Formats) + + marshalled, err := parsed.Marshal() + require.NoError(t, err) + require.Less(t, len(marshalled), len(offer.SDP)) + + // audio is shrunk only when included + parsed, err = offer.Unmarshal() + require.NoError(t, err) + require.Equal(t, 2, shrinkInactiveMediaSections(parsed, knownMids, true)) + require.Equal(t, []string{"111"}, parsed.MediaDescriptions[3].MediaName.Formats) + rtpmap, ok = parsed.MediaDescriptions[3].Attribute("rtpmap") + require.True(t, ok) + require.Equal(t, "111 opus/48000/2", rtpmap) +} + +// pion must accept a shrunk offer for a transceiver it already has and answer it as before +func TestShrinkInactiveMediaSectionsWithPion(t *testing.T) { + newTransport := func(shrink bool) *PCTransport { + tr, err := NewPCTransport(TransportParams{ + Config: &WebRTCConfig{}, + EnabledPublishCodecs: []*livekit.Codec{{Mime: mime.MimeTypeVP8.String()}}, + Handler: &transportfakes.FakeHandler{}, + ShrinkInactiveMediaSections: config.ShrinkInactiveMediaSectionsConfig{ + Enabled: shrink, + MinSDPSize: 1, + }, + }) + require.NoError(t, err) + return tr + } + negotiate := func(tr *PCTransport, sections []testOfferSection) string { + require.NoError(t, tr.setRemoteDescription(webrtc.SessionDescription{Type: webrtc.SDPTypeOffer, SDP: testOffer(sections)})) + answer, err := tr.pc.CreateAnswer(nil) + require.NoError(t, err) + require.NoError(t, tr.pc.SetLocalDescription(answer)) + return answer.SDP + } + // only direction and mid lines matter for the answer comparison, ICE candidates differ per run + answerShape := func(answer string) []string { + var shape []string + for _, line := range strings.Split(answer, "\r\n") { + if strings.HasPrefix(line, "m=") || strings.HasPrefix(line, "a=mid:") || line == "a=inactive" || line == "a=recvonly" || line == "a=sendonly" { + shape = append(shape, line) + } + } + return shape + } + + sequence := [][]testOfferSection{ + {{mid: 0}}, + {{mid: 0, inactive: true}, {mid: 1}}, + {{mid: 0, inactive: true}, {mid: 1, inactive: true}, {mid: 2}}, + } + + plain := newTransport(false) + defer plain.Close() + shrunk := newTransport(true) + defer shrunk.Close() + for i, sections := range sequence { + plainAnswer := negotiate(plain, sections) + shrunkAnswer := negotiate(shrunk, sections) + require.Equal(t, answerShape(plainAnswer), answerShape(shrunkAnswer), "offer %d", i) + } + + // the shrunk transport gave pion a smaller remote description + require.Less(t, len(shrunk.pc.CurrentRemoteDescription().SDP), len(plain.pc.CurrentRemoteDescription().SDP)) + require.Equal(t, 3, len(shrunk.pc.GetTransceivers())) + // pion kept the codecs negotiated when the sections were active + for _, tr := range shrunk.pc.GetTransceivers() { + require.NotEmpty(t, tr.Receiver().GetParameters().Codecs, "mid %s", tr.Mid()) + } +} diff --git a/pkg/rtc/transportmanager.go b/pkg/rtc/transportmanager.go index 108651f21..cf58bd120 100644 --- a/pkg/rtc/transportmanager.go +++ b/pkg/rtc/transportmanager.go @@ -100,6 +100,7 @@ type TransportManagerParams struct { FireOnTrackBySdp bool EnableDataTracks bool ExcludeIPv6LocalCandidates bool + ShrinkInactiveMediaSections config.ShrinkInactiveMediaSectionsConfig EnableWarp bool } @@ -173,6 +174,7 @@ func NewTransportManager(params TransportManagerParams) (*TransportManager, erro FireOnTrackBySdp: params.FireOnTrackBySdp, EnableDataTracks: params.EnableDataTracks, ExcludeIPv6LocalCandidates: params.ExcludeIPv6LocalCandidates, + ShrinkInactiveMediaSections: params.ShrinkInactiveMediaSections, EnableWarp: params.EnableWarp, }) if err != nil { diff --git a/pkg/service/roommanager.go b/pkg/service/roommanager.go index e25fae883..de1bbb98d 100644 --- a/pkg/service/roommanager.go +++ b/pkg/service/roommanager.go @@ -542,6 +542,7 @@ func (r *RoomManager) StartSession( EnableDataTracks: r.config.EnableDataTracks, EnableParticipantDataBlob: r.config.EnableParticipantDataBlob, EnableRTPStreamRestartDetection: r.config.RTC.EnableRTPStreamRestartDetection, + ShrinkInactiveMediaSections: r.config.RTC.ShrinkInactiveMediaSections, EnableWarp: enableWarp, }) if err != nil {