Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .changes/negotiation-error-reconnect
Original file line number Diff line number Diff line change
@@ -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`"
22 changes: 12 additions & 10 deletions lib/src/core/engine.dart
Original file line number Diff line number Diff line change
Expand Up @@ -344,17 +344,17 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
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<void> _onPublisherNegotiationError(Object error) async {
if (error is NegotiationError) {
fullReconnectOnNext = true;
}
await handleReconnect(
ClientDisconnectReason.negotiationFailed,
reconnectReason: lk_models.ReconnectReason.RR_UNKNOWN,
);
}

bool? isBufferStatusLow(Reliability kind) {
Expand Down Expand Up @@ -698,6 +698,8 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
signalClient.sendOffer(offer);
};

publisher?.onNegotiationError = _onPublisherNegotiationError;

// in subscriber primary mode, server side opens sub data channels.
if (_subscriberPrimary) {
subscriber?.pc.onDataChannel = _onDataChannel;
Expand Down
25 changes: 23 additions & 2 deletions lib/src/core/transport.dart
Original file line number Diff line number Diff line change
Expand Up @@ -172,6 +172,7 @@ void applyVideoStartBitrate(Map<String, dynamic> media, int codecPayload, int st
}

typedef TransportOnOffer = void Function(rtc.RTCSessionDescription offer);
typedef TransportOnNegotiationError = void Function(Object error);
typedef PeerConnectionCreate =
Future<rtc.RTCPeerConnection> Function(Map<String, dynamic> configuration, [Map<String, dynamic> constraints]);

Expand All @@ -191,6 +192,7 @@ class Transport extends Disposable {
bool restartingIce = false;
bool renegotiate = false;
TransportOnOffer? onOffer;
TransportOnNegotiationError? onNegotiationError;
Function? _cancelDebounce;
ConnectOptions connectOptions;

Expand Down Expand Up @@ -241,11 +243,23 @@ class Transport extends Disposable {
}

late final negotiate = Utils.createDebounceFunc(
(void _) => createAndSendOffer(),
(void _) => _createAndSendOfferReportingErrors(),
cancelFunc: (f) => _cancelDebounce = f,
wait: connectOptions.timeouts.debounce,
);

/// 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<void> _createAndSendOfferReportingErrors() async {
try {
await createAndSendOffer();
} catch (error) {
logger.warning('[$objectId] negotiate() failed with error: $error');
onNegotiationError?.call(error);
Comment thread
devin-ai-integration[bot] marked this conversation as resolved.
}
}

Future<void> setRemoteDescription(rtc.RTCSessionDescription sd) async {
if (isDisposed) {
logger.warning('[$objectId] setRemoteDescription() already disposed');
Expand All @@ -267,7 +281,14 @@ 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. Without a handler the
// error propagates to the caller instead of being dropped.
if (onNegotiationError == null) {
await createAndSendOffer();
} else {
await _createAndSendOfferReportingErrors();
}
}
}

Expand Down
140 changes: 140 additions & 0 deletions test/core/transport_negotiate_test.dart
Original file line number Diff line number Diff line change
@@ -0,0 +1,140 @@
// 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<void> setLocalDescription(rtc.RTCSessionDescription description) async {
throw Exception('The order of m-lines in subsequent offer doesn\'t match order from previous offer/answer.');
}
}

Future<rtc.RTCPeerConnection> _createRejecting(
Map<String, dynamic> configuration, [
Map<String, dynamic>? constraints,
]) async => _RejectingPeerConnection();

class _RejectAfterAnswerPeerConnection extends MockPeerConnection {
bool rejectLocal = false;

@override
Future<void> 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 {
final uncaught = <Object>[];
final reported = Completer<Object>();

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<void>.delayed(const Duration(milliseconds: 200));
}, (error, stack) => uncaught.add(error));

expect(uncaught, isEmpty);
expect(reported.isCompleted, isTrue);
expect(await reported.future, isA<NegotiationError>());
});

test('does not report when the offer is sent', () async {
final offers = <rtc.RTCSessionDescription>[];
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<void>.delayed(const Duration(milliseconds: 200));

expect(offers, hasLength(1));
});

test('reports a failed deferred offer to onNegotiationError', () async {
final pc = _RejectAfterAnswerPeerConnection();
final transport = await Transport.create(
(Map<String, dynamic> configuration, [Map<String, dynamic>? constraints]) async => pc,
connectOptions: const ConnectOptions(),
);
addTearDown(transport.dispose);
transport.onOffer = (_) {};
final reported = Completer<Object>();
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<NegotiationError>());
});

test('throws from setRemoteDescription when a deferred offer fails and no handler is set', () async {
final pc = _RejectAfterAnswerPeerConnection();
final transport = await Transport.create(
(Map<String, dynamic> configuration, [Map<String, dynamic>? 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<NegotiationError>()),
);
});

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<NegotiationError>()));
});
});
}
Loading