From 38edeac033c2465ef23f80b84db4d33145d88d43 Mon Sep 17 00:00:00 2001 From: Raphael Fakhri <153192858+RaphaelFakhri@users.noreply.github.com> Date: Tue, 29 Sep 2026 09:42:51 +0000 Subject: [PATCH 1/3] fix: route publisher negotiation errors to reconnect instead of leaking them --- .changes/negotiation-error-reconnect | 1 + lib/src/core/engine.dart | 22 +++---- lib/src/core/transport.dart | 15 ++++- test/core/transport_negotiate_test.dart | 76 +++++++++++++++++++++++++ 4 files changed, 103 insertions(+), 11 deletions(-) create mode 100644 .changes/negotiation-error-reconnect create mode 100644 test/core/transport_negotiate_test.dart diff --git a/.changes/negotiation-error-reconnect b/.changes/negotiation-error-reconnect new file mode 100644 index 000000000..c41fdd232 --- /dev/null +++ b/.changes/negotiation-error-reconnect @@ -0,0 +1 @@ +patch type="fixed" "A failed publisher negotiation (for example a `setLocalDescription` rejection) is now routed to the reconnect path instead of escaping as an unhandled `NegotiationError`" diff --git a/lib/src/core/engine.dart b/lib/src/core/engine.dart index ba82dfefb..4d7fa3e7f 100644 --- a/lib/src/core/engine.dart +++ b/lib/src/core/engine.dart @@ -344,17 +344,17 @@ class Engine extends Disposable with EventsEmittable { return; } _hasPublished = true; - try { - publisher!.negotiate(null); - } catch (error) { - if (error is NegotiationError) { - fullReconnectOnNext = true; - } - await handleReconnect( - ClientDisconnectReason.negotiationFailed, - reconnectReason: lk_models.ReconnectReason.RR_UNKNOWN, - ); + publisher!.negotiate(null); + } + + Future _onPublisherNegotiationError(Object error) async { + if (error is NegotiationError) { + fullReconnectOnNext = true; } + await handleReconnect( + ClientDisconnectReason.negotiationFailed, + reconnectReason: lk_models.ReconnectReason.RR_UNKNOWN, + ); } bool? isBufferStatusLow(Reliability kind) { @@ -698,6 +698,8 @@ class Engine extends Disposable with EventsEmittable { signalClient.sendOffer(offer); }; + publisher?.onNegotiationError = _onPublisherNegotiationError; + // in subscriber primary mode, server side opens sub data channels. if (_subscriberPrimary) { subscriber?.pc.onDataChannel = _onDataChannel; diff --git a/lib/src/core/transport.dart b/lib/src/core/transport.dart index a6a605900..32dc42b3c 100644 --- a/lib/src/core/transport.dart +++ b/lib/src/core/transport.dart @@ -172,6 +172,7 @@ void applyVideoStartBitrate(Map media, int codecPayload, int st } typedef TransportOnOffer = void Function(rtc.RTCSessionDescription offer); +typedef TransportOnNegotiationError = void Function(Object error); typedef PeerConnectionCreate = Future Function(Map configuration, [Map constraints]); @@ -191,6 +192,7 @@ class Transport extends Disposable { bool restartingIce = false; bool renegotiate = false; TransportOnOffer? onOffer; + TransportOnNegotiationError? onNegotiationError; Function? _cancelDebounce; ConnectOptions connectOptions; @@ -241,11 +243,22 @@ class Transport extends Disposable { } late final negotiate = Utils.createDebounceFunc( - (void _) => createAndSendOffer(), + (void _) => _createAndSendOfferReportingErrors(), cancelFunc: (f) => _cancelDebounce = f, wait: connectOptions.timeouts.debounce, ); + /// The debouncer discards the returned future, so a failure here would surface as an + /// unhandled error. Hand it to [onNegotiationError] instead. + Future _createAndSendOfferReportingErrors() async { + try { + await createAndSendOffer(); + } catch (error) { + logger.warning('[$objectId] negotiate() failed with error: $error'); + onNegotiationError?.call(error); + } + } + Future setRemoteDescription(rtc.RTCSessionDescription sd) async { if (isDisposed) { logger.warning('[$objectId] setRemoteDescription() already disposed'); diff --git a/test/core/transport_negotiate_test.dart b/test/core/transport_negotiate_test.dart new file mode 100644 index 000000000..a5e5b2fe8 --- /dev/null +++ b/test/core/transport_negotiate_test.dart @@ -0,0 +1,76 @@ +// 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. + +@Timeout(Duration(seconds: 5)) +library; + +import 'dart:async'; + +import 'package:flutter_test/flutter_test.dart'; +import 'package:flutter_webrtc/flutter_webrtc.dart' as rtc; + +import 'package:livekit_client/src/core/transport.dart' show Transport; +import 'package:livekit_client/src/exceptions.dart' show NegotiationError; +import 'package:livekit_client/src/options.dart' show ConnectOptions; +import '../mock/peerconnection_mock.dart'; + +class _RejectingPeerConnection extends MockPeerConnection { + @override + Future setLocalDescription(rtc.RTCSessionDescription description) async { + throw Exception('The order of m-lines in subsequent offer doesn\'t match order from previous offer/answer.'); + } +} + +Future _createRejecting( + Map configuration, [ + Map? constraints, +]) async => _RejectingPeerConnection(); + +void main() { + group('Transport.negotiate', () { + test('reports a failed offer to onNegotiationError instead of leaking an unhandled error', () async { + final uncaught = []; + final reported = Completer(); + + await runZonedGuarded(() async { + final transport = await Transport.create(_createRejecting, connectOptions: const ConnectOptions()); + addTearDown(transport.dispose); + transport.onOffer = (_) {}; + transport.onNegotiationError = reported.complete; + + transport.negotiate(null); + + // Outlast the debounce and the rejected setLocalDescription. + await Future.delayed(const Duration(milliseconds: 200)); + }, (error, stack) => uncaught.add(error)); + + expect(uncaught, isEmpty); + expect(reported.isCompleted, isTrue); + expect(await reported.future, isA()); + }); + + test('does not report when the offer is sent', () async { + final offers = []; + final transport = await Transport.create(MockPeerConnection.create, connectOptions: const ConnectOptions()); + addTearDown(transport.dispose); + transport.onOffer = offers.add; + transport.onNegotiationError = (error) => fail('unexpected negotiation error: $error'); + + transport.negotiate(null); + await Future.delayed(const Duration(milliseconds: 200)); + + expect(offers, hasLength(1)); + }); + }); +} From 9646ad1805cf2b2d04dca3ae5399bcf86ddbc31a Mon Sep 17 00:00:00 2001 From: Raphael Fakhri <153192858+RaphaelFakhri@users.noreply.github.com> Date: Tue, 29 Sep 2026 10:26:09 +0000 Subject: [PATCH 2/3] fix: report failed deferred publisher offers to reconnect --- lib/src/core/transport.dart | 9 +++-- test/core/transport_negotiate_test.dart | 44 +++++++++++++++++++++++++ 2 files changed, 50 insertions(+), 3 deletions(-) diff --git a/lib/src/core/transport.dart b/lib/src/core/transport.dart index 32dc42b3c..4aa2682dd 100644 --- a/lib/src/core/transport.dart +++ b/lib/src/core/transport.dart @@ -248,8 +248,9 @@ class Transport extends Disposable { wait: connectOptions.timeouts.debounce, ); - /// The debouncer discards the returned future, so a failure here would surface as an - /// unhandled error. Hand it to [onNegotiationError] instead. + /// The debouncer and the offer deferred until the answer arrives have no caller that handles + /// errors, so a failure here would surface as an unhandled error. Hand it to + /// [onNegotiationError] instead. Direct [createAndSendOffer] callers still get the error. Future _createAndSendOfferReportingErrors() async { try { await createAndSendOffer(); @@ -280,7 +281,9 @@ class Transport extends Disposable { if (renegotiate) { renegotiate = false; - await createAndSendOffer(); // await or un-awaited ? + // The signal listener that awaits this call has no reconnect handling, so a failed deferred + // offer is reported to [onNegotiationError] like a debounced one. + await _createAndSendOfferReportingErrors(); } } diff --git a/test/core/transport_negotiate_test.dart b/test/core/transport_negotiate_test.dart index a5e5b2fe8..7c0d6de75 100644 --- a/test/core/transport_negotiate_test.dart +++ b/test/core/transport_negotiate_test.dart @@ -37,6 +37,18 @@ Future _createRejecting( Map? constraints, ]) async => _RejectingPeerConnection(); +class _RejectAfterAnswerPeerConnection extends MockPeerConnection { + bool rejectLocal = false; + + @override + Future setLocalDescription(rtc.RTCSessionDescription description) async { + if (rejectLocal) { + throw Exception('The order of m-lines in subsequent offer doesn\'t match order from previous offer/answer.'); + } + await super.setLocalDescription(description); + } +} + void main() { group('Transport.negotiate', () { test('reports a failed offer to onNegotiationError instead of leaking an unhandled error', () async { @@ -72,5 +84,37 @@ void main() { expect(offers, hasLength(1)); }); + + test('reports a failed deferred offer to onNegotiationError', () async { + final pc = _RejectAfterAnswerPeerConnection(); + final transport = await Transport.create( + (Map configuration, [Map? constraints]) async => pc, + connectOptions: const ConnectOptions(), + ); + addTearDown(transport.dispose); + transport.onOffer = (_) {}; + final reported = Completer(); + transport.onNegotiationError = reported.complete; + + // An offer is already waiting for its answer, so the next one is deferred. + await pc.setLocalDescription(await pc.createOffer()); + await transport.createAndSendOffer(); + expect(transport.renegotiate, isTrue); + + pc.rejectLocal = true; + await transport.setRemoteDescription(rtc.RTCSessionDescription('v=0', 'answer')); + + expect(reported.isCompleted, isTrue); + expect(await reported.future, isA()); + }); + + test('still throws from a direct createAndSendOffer', () async { + final transport = await Transport.create(_createRejecting, connectOptions: const ConnectOptions()); + addTearDown(transport.dispose); + transport.onOffer = (_) {}; + transport.onNegotiationError = (error) => fail('unexpected negotiation error: $error'); + + await expectLater(transport.createAndSendOffer(), throwsA(isA())); + }); }); } From 1793abf13b2882cecb18686f10b48b7738f9490d Mon Sep 17 00:00:00 2001 From: Raphael Fakhri <153192858+RaphaelFakhri@users.noreply.github.com> Date: Wed, 30 Sep 2026 09:16:16 +0000 Subject: [PATCH 3/3] fix: propagate a failed deferred offer when no negotiation handler is set --- lib/src/core/transport.dart | 9 +++++++-- test/core/transport_negotiate_test.dart | 20 ++++++++++++++++++++ 2 files changed, 27 insertions(+), 2 deletions(-) diff --git a/lib/src/core/transport.dart b/lib/src/core/transport.dart index 4aa2682dd..084104336 100644 --- a/lib/src/core/transport.dart +++ b/lib/src/core/transport.dart @@ -282,8 +282,13 @@ class Transport extends Disposable { if (renegotiate) { renegotiate = false; // The signal listener that awaits this call has no reconnect handling, so a failed deferred - // offer is reported to [onNegotiationError] like a debounced one. - await _createAndSendOfferReportingErrors(); + // offer is reported to [onNegotiationError] like a debounced one. Without a handler the + // error propagates to the caller instead of being dropped. + if (onNegotiationError == null) { + await createAndSendOffer(); + } else { + await _createAndSendOfferReportingErrors(); + } } } diff --git a/test/core/transport_negotiate_test.dart b/test/core/transport_negotiate_test.dart index 7c0d6de75..d8f257a71 100644 --- a/test/core/transport_negotiate_test.dart +++ b/test/core/transport_negotiate_test.dart @@ -108,6 +108,26 @@ void main() { expect(await reported.future, isA()); }); + test('throws from setRemoteDescription when a deferred offer fails and no handler is set', () async { + final pc = _RejectAfterAnswerPeerConnection(); + final transport = await Transport.create( + (Map configuration, [Map? constraints]) async => pc, + connectOptions: const ConnectOptions(), + ); + addTearDown(transport.dispose); + transport.onOffer = (_) {}; + + await pc.setLocalDescription(await pc.createOffer()); + await transport.createAndSendOffer(); + expect(transport.renegotiate, isTrue); + + pc.rejectLocal = true; + await expectLater( + transport.setRemoteDescription(rtc.RTCSessionDescription('v=0', 'answer')), + throwsA(isA()), + ); + }); + test('still throws from a direct createAndSendOffer', () async { final transport = await Transport.create(_createRejecting, connectOptions: const ConnectOptions()); addTearDown(transport.dispose);