diff --git a/example/lib/exts.dart b/example/lib/exts.dart index ece86d9..15f4ab5 100644 --- a/example/lib/exts.dart +++ b/example/lib/exts.dart @@ -89,6 +89,20 @@ extension LKExampleExt on BuildContext { ), ); + Future showReconnectSuccessDialog() => showDialog( + 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 showSendDataDialog() => showDialog( context: this, builder: (ctx) => AlertDialog( diff --git a/example/lib/widgets/controls.dart b/example/lib/widgets/controls.dart index a1ebad3..f7c7064 100644 --- a/example/lib/widgets/controls.dart +++ b/example/lib/widgets/controls.dart @@ -126,7 +126,14 @@ class _ControlsWidgetState extends State { void _onTapReconnect() async { 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 { diff --git a/lib/src/core/engine.dart b/lib/src/core/engine.dart index 4286600..6401d66 100644 --- a/lib/src/core/engine.dart +++ b/lib/src/core/engine.dart @@ -86,8 +86,7 @@ class Engine extends Disposable with EventsEmittable { }) : signalClient = signalClient ?? SignalClient() { if (kDebugMode) { // log all EngineEvents - events.listen((event) => - logger.fine('[EngineEvent] $objectId ${event.runtimeType}')); + events.listen((event) => logger.fine('[EngineEvent] $objectId ${event}')); } _setUpListeners(); @@ -241,7 +240,7 @@ class Engine extends Disposable with EventsEmittable { @internal Future reconnect() async { if (_connectionState == ConnectionState.disconnected) { - logger.fine('$objectId reconnect() already closed'); + logger.fine('Reconnect: Already closed.'); return; } @@ -252,6 +251,8 @@ class Engine extends Disposable with EventsEmittable { throw ConnectException('could not reconnect without url and token'); } + _connectionState = ConnectionState.reconnecting; + if (_reconnectAttempts == 0) { events.emit(const EngineReconnectingEvent()); } @@ -259,7 +260,6 @@ class Engine extends Disposable with EventsEmittable { try { // isReconnecting = true; - _connectionState = ConnectionState.reconnecting; await signalClient.reconnect( url, token, @@ -274,13 +274,19 @@ class Engine extends Disposable with EventsEmittable { // await negotiate(iceRestart: true); if (_hasPublished) { - logger.fine('reconnect: publisher.createAndSendOffer'); - await publisher! - .createAndSendOffer(const RTCOfferOptions(iceRestart: true)); + logger.fine('Reconnect: negotiating publisher...'); + await publisher!.createAndSendOffer(const RTCOfferOptions( + iceRestart: true, + )); } - if (!(primary?.pc.iceConnectionState?.isConnected() ?? false)) { - logger.fine('reconnect: waiting for primary to ice-connect...'); + final iceConnected = + primary?.pc.iceConnectionState?.isConnected() ?? false; + + logger.fine('Reconnect: iceConnected: $iceConnected'); + + if (!iceConnected) { + logger.fine('Reconnect: Waiting for primary to connect...'); await events.waitFor( filter: (event) => event.isPrimary && event.iceState.isConnected(), @@ -288,15 +294,15 @@ class Engine extends Disposable with EventsEmittable { ); } - logger.fine('reconnect: success'); + logger.fine('Reconnect: success'); + _connectionState = ConnectionState.connected; events.emit(const EngineReconnectedEvent()); _reconnectAttempts = 0; - - // don't catch and pass up any exception - } finally { - // always set reconnecting to false - // isReconnecting = false; + } catch (error) { + logger.fine('Reconnect: error ${error}'); + // Pass up all exceptions _connectionState = ConnectionState.disconnected; + rethrow; } } @@ -344,7 +350,7 @@ class Engine extends Disposable with EventsEmittable { subscriber?.pc.onIceConnectionState = (state) => events.emit(EngineSubscriberIceStateUpdatedEvent( - state: state, + iceState: state, isPrimary: _subscriberPrimary, )); @@ -491,6 +497,7 @@ class Engine extends Disposable with EventsEmittable { } Future _onDisconnected(String reason) async { + logger.info('onDisconnected reason: $reason'); if (_connectionState == ConnectionState.disconnected) { logger.fine('[$objectId] Already disconnected $reason'); return; @@ -596,6 +603,10 @@ class Engine extends Disposable with EventsEmittable { // relay ..on((event) => events.emit(event)) ..on((event) async { + if (connectionState == ConnectionState.reconnecting) { + logger.fine('Ignoring leave signal since engine is reconnecting...'); + return; + } await close(); events.emit(const EngineDisconnectedEvent()); }) diff --git a/lib/src/core/room.dart b/lib/src/core/room.dart index 29facab..9a95945 100644 --- a/lib/src/core/room.dart +++ b/lib/src/core/room.dart @@ -18,6 +18,7 @@ import '../proto/livekit_rtc.pb.dart' as lk_rtc; import '../support/disposable.dart'; import '../track/track.dart'; import '../types.dart'; +import '../core/signal_client.dart'; import 'engine.dart'; /// Room is the primary construct for LiveKit conferences. It contains a diff --git a/lib/src/core/signal_client.dart b/lib/src/core/signal_client.dart index a9bc960..3060a97 100644 --- a/lib/src/core/signal_client.dart +++ b/lib/src/core/signal_client.dart @@ -18,8 +18,9 @@ import '../types.dart'; import '../utils.dart'; class SignalClient extends Disposable with EventsEmittable { - // - bool _connected = false; + // Connection state of the socket conection. + ConnectionState _connectionState = ConnectionState.disconnected; + LiveKitWebSocket? _ws; SignalClient() { @@ -33,8 +34,6 @@ class SignalClient extends Disposable with EventsEmittable { }); } - bool get connected => _connected; - Future connect( String uriString, String token, { @@ -51,7 +50,7 @@ class SignalClient extends Disposable with EventsEmittable { rtcUri, WebSocketEventHandlers( onData: _onSocketData, - onDispose: _onSocketDone, + onDispose: _onSocketDispose, onError: _handleError, ), ); @@ -86,7 +85,8 @@ class SignalClient extends Disposable with EventsEmittable { String token, { ConnectOptions? connectOptions, }) async { - _connected = false; + logger.fine('SignalClient reconnecting...'); + _connectionState = ConnectionState.reconnecting; await _ws?.dispose(); _ws = null; @@ -101,19 +101,106 @@ class SignalClient extends Disposable with EventsEmittable { rtcUri, WebSocketEventHandlers( onData: _onSocketData, - onDispose: _onSocketDone, + onDispose: _onSocketDispose, onError: _handleError, ), ); - _connected = true; + logger.fine('SignalClient socket reconnected'); + _connectionState = ConnectionState.connected; } Future close() async { - _connected = false; + logger.fine('SignalClient close'); 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 _onSocketData(dynamic message) async { + if (message is! List) 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) => _sendRequest(lk_rtc.SignalRequest( offer: offer.toSDKType(), @@ -206,87 +293,4 @@ class SignalClient extends Disposable with EventsEmittable { void sendLeave() => _sendRequest(lk_rtc.SignalRequest( 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 _onSocketData(dynamic message) async { - if (message is! List) 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()); - } } diff --git a/lib/src/internal/events.dart b/lib/src/internal/events.dart index cd6cf0b..7852f98 100644 --- a/lib/src/internal/events.dart +++ b/lib/src/internal/events.dart @@ -23,12 +23,16 @@ abstract class EngineIceStateUpdatedEvent with EngineEvent, InternalEvent { @internal class EngineSubscriberIceStateUpdatedEvent extends EngineIceStateUpdatedEvent { const EngineSubscriberIceStateUpdatedEvent({ - required rtc.RTCIceConnectionState state, + required rtc.RTCIceConnectionState iceState, required bool isPrimary, }) : super( - iceState: state, + iceState: iceState, isPrimary: isPrimary, ); + + @override + String toString() => + '${runtimeType}(state: ${iceState}, isPrimary: ${isPrimary})'; } @internal @@ -40,6 +44,9 @@ class EnginePublisherIceStateUpdatedEvent extends EngineIceStateUpdatedEvent { iceState: state, isPrimary: isPrimary, ); + @override + String toString() => + '${runtimeType}(state: ${iceState}, isPrimary: ${isPrimary})'; } @internal diff --git a/lib/src/participant/local.dart b/lib/src/participant/local.dart index 7ca02ac..2899870 100644 --- a/lib/src/participant/local.dart +++ b/lib/src/participant/local.dart @@ -1,7 +1,5 @@ import 'package:flutter/foundation.dart'; - import 'package:flutter_webrtc/flutter_webrtc.dart' as rtc; -import 'package:meta/meta.dart'; import '../core/room.dart'; import '../events.dart'; diff --git a/lib/src/publication/remote.dart b/lib/src/publication/remote.dart index f6126b6..5e67db1 100644 --- a/lib/src/publication/remote.dart +++ b/lib/src/publication/remote.dart @@ -5,6 +5,7 @@ import 'package:collection/collection.dart'; import 'package:meta/meta.dart'; import '../events.dart'; +import '../core/signal_client.dart'; import '../extensions.dart'; import '../internal/events.dart'; import '../logger.dart'; diff --git a/lib/src/publication/track_publication.dart b/lib/src/publication/track_publication.dart index f5c9205..176daa1 100644 --- a/lib/src/publication/track_publication.dart +++ b/lib/src/publication/track_publication.dart @@ -9,6 +9,7 @@ import '../proto/livekit_models.pb.dart' as lk_models; import '../support/disposable.dart'; import '../track/track.dart'; import '../types.dart'; +import '../core/signal_client.dart'; /// Represents a track that's published to the server. This class contains /// metadata associated with tracks.