diff --git a/.changes/simulcast-codecs-lifecycle b/.changes/simulcast-codecs-lifecycle new file mode 100644 index 000000000..5b4972408 --- /dev/null +++ b/.changes/simulcast-codecs-lifecycle @@ -0,0 +1 @@ +patch type="fixed" "Backup codec state is cleared on unpublish and full reconnect, so republishing no longer acts on senders from a torn down connection" diff --git a/lib/src/participant/local.dart b/lib/src/participant/local.dart index b367f6c9e..51b62111a 100644 --- a/lib/src/participant/local.dart +++ b/lib/src/participant/local.dart @@ -552,23 +552,41 @@ class LocalParticipant extends Participant { } final sender = track.transceiver?.sender; + var didRemoveSender = false; if (sender != null) { try { await room.engine.publisher?.pc.removeTrack(sender); - if (track is LocalVideoTrack) { - track.simulcastCodecs.forEach((key, simulcastTrack) async { - await room.engine.publisher?.pc.removeTrack(simulcastTrack.sender!); - }); - } } catch (e) { logger.warning('[$objectId] rtc.removeTrack() did throw $e'); } + didRemoveSender = true; + } - // doesn't make sense to negotiate if already disposed - if (!isDisposed) { - // manual negotiation since track changed - await room.engine.negotiate(); + // not gated on the primary sender, stale backup codec state must not + // survive unpublish even when the track never got a live sender + if (track is LocalVideoTrack) { + // remove each backup sender on its own, one failure should not + // prevent removal of the others + for (final simulcastTrack in track.simulcastCodecs.values.toList()) { + final simulcastSender = simulcastTrack.sender; + if (simulcastSender == null) { + continue; + } + try { + await room.engine.publisher?.pc.removeTrack(simulcastSender); + } catch (e) { + logger.warning('[$objectId] rtc.removeTrack() did throw $e'); + } + simulcastTrack.sender = null; + didRemoveSender = true; } + track.clearSimulcastState(); + } + + // doesn't make sense to negotiate if already disposed + if (didRemoveSender && !isDisposed) { + // manual negotiation since track changed + await room.engine.negotiate(); } // did unpublish @@ -605,7 +623,11 @@ class LocalParticipant extends Participant { if (track.track is LocalAudioTrack) { await publishAudioTrack(track.track as LocalAudioTrack); } else if (track.track is LocalVideoTrack) { - await publishVideoTrack(track.track as LocalVideoTrack); + final videoTrack = track.track as LocalVideoTrack; + // a full reconnect replaced the peer connection, so any simulcast + // codec senders the track still holds belong to the old one + videoTrack.clearSimulcastState(); + await publishVideoTrack(videoTrack); } } } diff --git a/lib/src/track/local/video.dart b/lib/src/track/local/video.dart index c6b0c8239..d1e4773dc 100644 --- a/lib/src/track/local/video.dart +++ b/lib/src/track/local/video.dart @@ -504,6 +504,18 @@ extension LocalVideoTrackExt on LocalVideoTrack { return simulcastCodecInfo; } + /// Drops all simulcast codec state tied to the current publish session. + /// + /// Must be called when the track's senders are no longer valid, on unpublish + /// and before republishing after a full reconnect. Otherwise later publishes + /// see stale senders and [addSimulcastTrack] rejects the codec as a duplicate + /// when the server requests the backup codec again. + @internal + void clearSimulcastState() { + simulcastCodecs.clear(); + encodingBackups.clear(); + } + Future setDegradationPreference(DegradationPreference preference) async { _degradationPreference = preference; await applyDegradationPreference(sender); diff --git a/test/track/simulcast_state_test.dart b/test/track/simulcast_state_test.dart new file mode 100644 index 000000000..26992194c --- /dev/null +++ b/test/track/simulcast_state_test.dart @@ -0,0 +1,165 @@ +// 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. + +import 'dart:typed_data'; + +import 'package:flutter_test/flutter_test.dart'; +import 'package:flutter_webrtc/flutter_webrtc.dart' as rtc; + +import 'package:livekit_client/src/track/local/video.dart'; +import 'package:livekit_client/src/track/options.dart'; +import 'package:livekit_client/src/types/other.dart'; + +void main() { + TestWidgetsFlutterBinding.ensureInitialized(); + + LocalVideoTrack createTrack() { + final mediaTrack = _FakeMediaStreamTrack(id: 'video-1', kind: 'video'); + final stream = _FakeMediaStream('stream-1'); + return LocalVideoTrack( + TrackSource.camera, + stream, + mediaTrack, + const CameraCaptureOptions(), + ); + } + + group('clearSimulcastState', () { + test('allows re-adding a backup codec after clearing', () { + final track = createTrack(); + + track.addSimulcastTrack('vp8', []); + // simulates the server requesting the same backup codec again after a + // reconnect, which previously threw because the map was never cleared + expect(() => track.addSimulcastTrack('vp8', []), throwsException); + + track.clearSimulcastState(); + expect(track.simulcastCodecs, isEmpty); + expect(() => track.addSimulcastTrack('vp8', []), returnsNormally); + }); + + test('clears encoding backups as well', () { + final track = createTrack(); + track.encodingBackups[('sender-1', 0)] = rtc.RTCRtpEncoding(); + + track.clearSimulcastState(); + + expect(track.encodingBackups, isEmpty); + }); + }); +} + +class _FakeMediaStream extends rtc.MediaStream { + final List _tracks = []; + + _FakeMediaStream(String id) : super(id, 'fake-owner'); + + @override + bool? get active => true; + + @override + Future addTrack(rtc.MediaStreamTrack track, {bool addToNative = true}) async { + _tracks.add(track); + } + + @override + Future clone() async => _FakeMediaStream('${id}_clone'); + + @override + List getAudioTracks() => _tracks.where((t) => t.kind == 'audio').toList(); + + @override + Future getMediaTracks() async {} + + @override + List getTracks() => List.from(_tracks); + + @override + List getVideoTracks() => _tracks.where((t) => t.kind == 'video').toList(); + + @override + Future removeTrack(rtc.MediaStreamTrack track, {bool removeFromNative = true}) async { + _tracks.remove(track); + } +} + +class _FakeMediaStreamTrack implements rtc.MediaStreamTrack { + @override + rtc.StreamTrackCallback? onEnded; + + @override + rtc.StreamTrackCallback? onMute; + + @override + rtc.StreamTrackCallback? onUnMute; + + @override + bool enabled; + + @override + final String id; + + @override + final String kind; + + @override + String? get label => '$kind-track'; + + @override + bool? get muted => false; + + _FakeMediaStreamTrack({ + required this.id, + required this.kind, + this.enabled = true, + }); + + @override + Future adaptRes(int width, int height) async {} + + @override + Future applyConstraints([Map? constraints]) async {} + + @override + Future captureFrame() { + throw UnimplementedError(); + } + + @override + Future clone() async => _FakeMediaStreamTrack(id: id, kind: kind, enabled: enabled); + + @override + Future dispose() async {} + + @override + Map getConstraints() => const {}; + + @override + Map getSettings() => const {}; + + @override + Future hasTorch() async => false; + + @override + void enableSpeakerphone(bool enable) {} + + @override + Future setTorch(bool torch) async {} + + @override + Future stop() async {} + + @override + Future switchCamera() async => false; +}