fix reconnect logic

This commit is contained in:
Hiroshi Horie
2021-12-23 01:50:58 +07:00
parent a0c7744baa
commit 94b4b07e92
9 changed files with 157 additions and 113 deletions
+14
View File
@@ -89,6 +89,20 @@ extension LKExampleExt on BuildContext {
), ),
); );
Future<void> showReconnectSuccessDialog() => showDialog<void>(
context: this,
builder: (ctx) => AlertDialog(
title: const Text('Reconnect'),
content: const Text('Reconnection was successful.'),
actions: [
TextButton(
onPressed: () => Navigator.pop(ctx),
child: const Text('OK'),
),
],
),
);
Future<bool?> showSendDataDialog() => showDialog<bool>( Future<bool?> showSendDataDialog() => showDialog<bool>(
context: this, context: this,
builder: (ctx) => AlertDialog( builder: (ctx) => AlertDialog(
+8 -1
View File
@@ -126,7 +126,14 @@ class _ControlsWidgetState extends State<ControlsWidget> {
void _onTapReconnect() async { void _onTapReconnect() async {
final result = await context.showReconnectDialog(); final result = await context.showReconnectDialog();
if (result == true) await widget.room.reconnect(); if (result == true) {
try {
await widget.room.reconnect();
await context.showReconnectSuccessDialog();
} catch (error) {
await context.showErrorDialog(error);
}
}
} }
void _onTapSendData() async { void _onTapSendData() async {
+27 -16
View File
@@ -86,8 +86,7 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
}) : signalClient = signalClient ?? SignalClient() { }) : signalClient = signalClient ?? SignalClient() {
if (kDebugMode) { if (kDebugMode) {
// log all EngineEvents // log all EngineEvents
events.listen((event) => events.listen((event) => logger.fine('[EngineEvent] $objectId ${event}'));
logger.fine('[EngineEvent] $objectId ${event.runtimeType}'));
} }
_setUpListeners(); _setUpListeners();
@@ -241,7 +240,7 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
@internal @internal
Future<void> reconnect() async { Future<void> reconnect() async {
if (_connectionState == ConnectionState.disconnected) { if (_connectionState == ConnectionState.disconnected) {
logger.fine('$objectId reconnect() already closed'); logger.fine('Reconnect: Already closed.');
return; return;
} }
@@ -252,6 +251,8 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
throw ConnectException('could not reconnect without url and token'); throw ConnectException('could not reconnect without url and token');
} }
_connectionState = ConnectionState.reconnecting;
if (_reconnectAttempts == 0) { if (_reconnectAttempts == 0) {
events.emit(const EngineReconnectingEvent()); events.emit(const EngineReconnectingEvent());
} }
@@ -259,7 +260,6 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
try { try {
// isReconnecting = true; // isReconnecting = true;
_connectionState = ConnectionState.reconnecting;
await signalClient.reconnect( await signalClient.reconnect(
url, url,
token, token,
@@ -274,13 +274,19 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
// await negotiate(iceRestart: true); // await negotiate(iceRestart: true);
if (_hasPublished) { if (_hasPublished) {
logger.fine('reconnect: publisher.createAndSendOffer'); logger.fine('Reconnect: negotiating publisher...');
await publisher! await publisher!.createAndSendOffer(const RTCOfferOptions(
.createAndSendOffer(const RTCOfferOptions(iceRestart: true)); iceRestart: true,
));
} }
if (!(primary?.pc.iceConnectionState?.isConnected() ?? false)) { final iceConnected =
logger.fine('reconnect: waiting for primary to ice-connect...'); primary?.pc.iceConnectionState?.isConnected() ?? false;
logger.fine('Reconnect: iceConnected: $iceConnected');
if (!iceConnected) {
logger.fine('Reconnect: Waiting for primary to connect...');
await events.waitFor<EngineIceStateUpdatedEvent>( await events.waitFor<EngineIceStateUpdatedEvent>(
filter: (event) => event.isPrimary && event.iceState.isConnected(), filter: (event) => event.isPrimary && event.iceState.isConnected(),
@@ -288,15 +294,15 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
); );
} }
logger.fine('reconnect: success'); logger.fine('Reconnect: success');
_connectionState = ConnectionState.connected;
events.emit(const EngineReconnectedEvent()); events.emit(const EngineReconnectedEvent());
_reconnectAttempts = 0; _reconnectAttempts = 0;
} catch (error) {
// don't catch and pass up any exception logger.fine('Reconnect: error ${error}');
} finally { // Pass up all exceptions
// always set reconnecting to false
// isReconnecting = false;
_connectionState = ConnectionState.disconnected; _connectionState = ConnectionState.disconnected;
rethrow;
} }
} }
@@ -344,7 +350,7 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
subscriber?.pc.onIceConnectionState = subscriber?.pc.onIceConnectionState =
(state) => events.emit(EngineSubscriberIceStateUpdatedEvent( (state) => events.emit(EngineSubscriberIceStateUpdatedEvent(
state: state, iceState: state,
isPrimary: _subscriberPrimary, isPrimary: _subscriberPrimary,
)); ));
@@ -491,6 +497,7 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
} }
Future<void> _onDisconnected(String reason) async { Future<void> _onDisconnected(String reason) async {
logger.info('onDisconnected reason: $reason');
if (_connectionState == ConnectionState.disconnected) { if (_connectionState == ConnectionState.disconnected) {
logger.fine('[$objectId] Already disconnected $reason'); logger.fine('[$objectId] Already disconnected $reason');
return; return;
@@ -596,6 +603,10 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
// relay // relay
..on<SignalStreamStateUpdatedEvent>((event) => events.emit(event)) ..on<SignalStreamStateUpdatedEvent>((event) => events.emit(event))
..on<SignalLeaveEvent>((event) async { ..on<SignalLeaveEvent>((event) async {
if (connectionState == ConnectionState.reconnecting) {
logger.fine('Ignoring leave signal since engine is reconnecting...');
return;
}
await close(); await close();
events.emit(const EngineDisconnectedEvent()); events.emit(const EngineDisconnectedEvent());
}) })
+1
View File
@@ -18,6 +18,7 @@ import '../proto/livekit_rtc.pb.dart' as lk_rtc;
import '../support/disposable.dart'; import '../support/disposable.dart';
import '../track/track.dart'; import '../track/track.dart';
import '../types.dart'; import '../types.dart';
import '../core/signal_client.dart';
import 'engine.dart'; import 'engine.dart';
/// Room is the primary construct for LiveKit conferences. It contains a /// Room is the primary construct for LiveKit conferences. It contains a
+96 -92
View File
@@ -18,8 +18,9 @@ import '../types.dart';
import '../utils.dart'; import '../utils.dart';
class SignalClient extends Disposable with EventsEmittable<SignalEvent> { class SignalClient extends Disposable with EventsEmittable<SignalEvent> {
// // Connection state of the socket conection.
bool _connected = false; ConnectionState _connectionState = ConnectionState.disconnected;
LiveKitWebSocket? _ws; LiveKitWebSocket? _ws;
SignalClient() { SignalClient() {
@@ -33,8 +34,6 @@ class SignalClient extends Disposable with EventsEmittable<SignalEvent> {
}); });
} }
bool get connected => _connected;
Future<void> connect( Future<void> connect(
String uriString, String uriString,
String token, { String token, {
@@ -51,7 +50,7 @@ class SignalClient extends Disposable with EventsEmittable<SignalEvent> {
rtcUri, rtcUri,
WebSocketEventHandlers( WebSocketEventHandlers(
onData: _onSocketData, onData: _onSocketData,
onDispose: _onSocketDone, onDispose: _onSocketDispose,
onError: _handleError, onError: _handleError,
), ),
); );
@@ -86,7 +85,8 @@ class SignalClient extends Disposable with EventsEmittable<SignalEvent> {
String token, { String token, {
ConnectOptions? connectOptions, ConnectOptions? connectOptions,
}) async { }) async {
_connected = false; logger.fine('SignalClient reconnecting...');
_connectionState = ConnectionState.reconnecting;
await _ws?.dispose(); await _ws?.dispose();
_ws = null; _ws = null;
@@ -101,19 +101,106 @@ class SignalClient extends Disposable with EventsEmittable<SignalEvent> {
rtcUri, rtcUri,
WebSocketEventHandlers( WebSocketEventHandlers(
onData: _onSocketData, onData: _onSocketData,
onDispose: _onSocketDone, onDispose: _onSocketDispose,
onError: _handleError, onError: _handleError,
), ),
); );
_connected = true; logger.fine('SignalClient socket reconnected');
_connectionState = ConnectionState.connected;
} }
Future<void> close() async { Future<void> close() async {
_connected = false; logger.fine('SignalClient close');
await _ws?.dispose(); await _ws?.dispose();
_ws = null;
} }
void _sendRequest(lk_rtc.SignalRequest req) {
if (_ws == null || isDisposed) {
logger.warning(
'[$objectId] Could not send message, not connected or already disposed');
return;
}
final buf = req.writeToBuffer();
_ws?.send(buf);
}
Future<void> _onSocketData(dynamic message) async {
if (message is! List<int>) return;
final msg = lk_rtc.SignalResponse.fromBuffer(message);
switch (msg.whichMessage()) {
case lk_rtc.SignalResponse_Message.join:
events.emit(SignalConnectedEvent(response: msg.join));
break;
case lk_rtc.SignalResponse_Message.answer:
events.emit(SignalAnswerEvent(sd: msg.answer.toSDKType()));
break;
case lk_rtc.SignalResponse_Message.offer:
events.emit(SignalOfferEvent(sd: msg.offer.toSDKType()));
break;
case lk_rtc.SignalResponse_Message.trickle:
events.emit(SignalTrickleEvent(
candidate: RTCIceCandidateExt.fromJson(msg.trickle.candidateInit),
target: msg.trickle.target,
));
break;
case lk_rtc.SignalResponse_Message.update:
events.emit(SignalParticipantUpdateEvent(
participants: msg.update.participants));
break;
case lk_rtc.SignalResponse_Message.trackPublished:
events.emit(SignalLocalTrackPublishedEvent(
cid: msg.trackPublished.cid,
track: msg.trackPublished.track,
));
break;
case lk_rtc.SignalResponse_Message.speakersChanged:
events.emit(
SignalSpeakersChangedEvent(speakers: msg.speakersChanged.speakers));
break;
case lk_rtc.SignalResponse_Message.connectionQuality:
events.emit(SignalConnectionQualityUpdateEvent(
updates: msg.connectionQuality.updates,
));
break;
case lk_rtc.SignalResponse_Message.leave:
events.emit(SignalLeaveEvent(canReconnect: msg.leave.canReconnect));
break;
case lk_rtc.SignalResponse_Message.mute:
events.emit(SignalMuteTrackEvent(
sid: msg.mute.sid,
muted: msg.mute.muted,
));
break;
case lk_rtc.SignalResponse_Message.streamStateUpdate:
events.emit(SignalStreamStateUpdatedEvent(
updates: msg.streamStateUpdate.streamStates,
));
break;
default:
logger.warning('skipping unsupported signal message');
}
}
void _handleError(dynamic error) {
logger.warning('received websocket error $error');
}
void _onSocketDispose() {
logger.fine('SignalClient onSocketDispose $_connectionState');
// don't emit event's when reconnecting state
if (_connectionState != ConnectionState.reconnecting) {
logger.fine('SignalClient did disconnect ${_connectionState}');
_connectionState = ConnectionState.disconnected;
events.emit(const SignalCloseEvent());
}
}
}
extension SignalClientRequests on SignalClient {
void sendOffer(rtc.RTCSessionDescription offer) => void sendOffer(rtc.RTCSessionDescription offer) =>
_sendRequest(lk_rtc.SignalRequest( _sendRequest(lk_rtc.SignalRequest(
offer: offer.toSDKType(), offer: offer.toSDKType(),
@@ -206,87 +293,4 @@ class SignalClient extends Disposable with EventsEmittable<SignalEvent> {
void sendLeave() => _sendRequest(lk_rtc.SignalRequest( void sendLeave() => _sendRequest(lk_rtc.SignalRequest(
leave: lk_rtc.LeaveRequest(), leave: lk_rtc.LeaveRequest(),
)); ));
void _sendRequest(lk_rtc.SignalRequest req) {
if (_ws == null || isDisposed) {
logger.warning(
'[$objectId] Could not send message, not connected or already disposed');
return;
}
final buf = req.writeToBuffer();
_ws?.send(buf);
}
Future<void> _onSocketData(dynamic message) async {
if (message is! List<int>) return;
final msg = lk_rtc.SignalResponse.fromBuffer(message);
switch (msg.whichMessage()) {
case lk_rtc.SignalResponse_Message.join:
if (!_connected) {
_connected = true;
events.emit(SignalConnectedEvent(response: msg.join));
}
break;
case lk_rtc.SignalResponse_Message.answer:
events.emit(SignalAnswerEvent(sd: msg.answer.toSDKType()));
break;
case lk_rtc.SignalResponse_Message.offer:
events.emit(SignalOfferEvent(sd: msg.offer.toSDKType()));
break;
case lk_rtc.SignalResponse_Message.trickle:
events.emit(SignalTrickleEvent(
candidate: RTCIceCandidateExt.fromJson(msg.trickle.candidateInit),
target: msg.trickle.target,
));
break;
case lk_rtc.SignalResponse_Message.update:
events.emit(SignalParticipantUpdateEvent(
participants: msg.update.participants));
break;
case lk_rtc.SignalResponse_Message.trackPublished:
events.emit(SignalLocalTrackPublishedEvent(
cid: msg.trackPublished.cid,
track: msg.trackPublished.track,
));
break;
case lk_rtc.SignalResponse_Message.speakersChanged:
events.emit(
SignalSpeakersChangedEvent(speakers: msg.speakersChanged.speakers));
break;
case lk_rtc.SignalResponse_Message.connectionQuality:
events.emit(SignalConnectionQualityUpdateEvent(
updates: msg.connectionQuality.updates,
));
break;
case lk_rtc.SignalResponse_Message.leave:
events.emit(SignalLeaveEvent(canReconnect: msg.leave.canReconnect));
break;
case lk_rtc.SignalResponse_Message.mute:
events.emit(SignalMuteTrackEvent(
sid: msg.mute.sid,
muted: msg.mute.muted,
));
break;
case lk_rtc.SignalResponse_Message.streamStateUpdate:
events.emit(SignalStreamStateUpdatedEvent(
updates: msg.streamStateUpdate.streamStates,
));
break;
default:
logger.warning('skipping unsupported signal message');
}
}
void _handleError(dynamic error) {
logger.warning('received websocket error $error');
}
void _onSocketDone() {
if (!_connected) return;
_ws = null;
_connected = false;
events.emit(const SignalCloseEvent());
}
} }
+9 -2
View File
@@ -23,12 +23,16 @@ abstract class EngineIceStateUpdatedEvent with EngineEvent, InternalEvent {
@internal @internal
class EngineSubscriberIceStateUpdatedEvent extends EngineIceStateUpdatedEvent { class EngineSubscriberIceStateUpdatedEvent extends EngineIceStateUpdatedEvent {
const EngineSubscriberIceStateUpdatedEvent({ const EngineSubscriberIceStateUpdatedEvent({
required rtc.RTCIceConnectionState state, required rtc.RTCIceConnectionState iceState,
required bool isPrimary, required bool isPrimary,
}) : super( }) : super(
iceState: state, iceState: iceState,
isPrimary: isPrimary, isPrimary: isPrimary,
); );
@override
String toString() =>
'${runtimeType}(state: ${iceState}, isPrimary: ${isPrimary})';
} }
@internal @internal
@@ -40,6 +44,9 @@ class EnginePublisherIceStateUpdatedEvent extends EngineIceStateUpdatedEvent {
iceState: state, iceState: state,
isPrimary: isPrimary, isPrimary: isPrimary,
); );
@override
String toString() =>
'${runtimeType}(state: ${iceState}, isPrimary: ${isPrimary})';
} }
@internal @internal
-2
View File
@@ -1,7 +1,5 @@
import 'package:flutter/foundation.dart'; import 'package:flutter/foundation.dart';
import 'package:flutter_webrtc/flutter_webrtc.dart' as rtc; import 'package:flutter_webrtc/flutter_webrtc.dart' as rtc;
import 'package:meta/meta.dart';
import '../core/room.dart'; import '../core/room.dart';
import '../events.dart'; import '../events.dart';
+1
View File
@@ -5,6 +5,7 @@ import 'package:collection/collection.dart';
import 'package:meta/meta.dart'; import 'package:meta/meta.dart';
import '../events.dart'; import '../events.dart';
import '../core/signal_client.dart';
import '../extensions.dart'; import '../extensions.dart';
import '../internal/events.dart'; import '../internal/events.dart';
import '../logger.dart'; import '../logger.dart';
@@ -9,6 +9,7 @@ import '../proto/livekit_models.pb.dart' as lk_models;
import '../support/disposable.dart'; import '../support/disposable.dart';
import '../track/track.dart'; import '../track/track.dart';
import '../types.dart'; import '../types.dart';
import '../core/signal_client.dart';
/// Represents a track that's published to the server. This class contains /// Represents a track that's published to the server. This class contains
/// metadata associated with tracks. /// metadata associated with tracks.