diff --git a/lib/src/core/engine.dart b/lib/src/core/engine.dart index 79e99e2..918c995 100644 --- a/lib/src/core/engine.dart +++ b/lib/src/core/engine.dart @@ -545,12 +545,15 @@ class Engine extends Disposable with EventsEmittable { if (_connectionState == ConnectionState.connected) { if (didReconnect) { events.emit(const EngineReconnectedEvent()); + // send queued requests if engine re-connected + signalClient.sendQueuedRequests(); } else { events.emit(const EngineConnectedEvent()); } } else if (_connectionState == ConnectionState.reconnecting) { events.emit(const EngineReconnectingEvent()); } else if (_connectionState == ConnectionState.disconnected) { + signalClient.cleanUp(); events.emit(const EngineDisconnectedEvent()); } } diff --git a/lib/src/core/signal_client.dart b/lib/src/core/signal_client.dart index f888892..ae3475f 100644 --- a/lib/src/core/signal_client.dart +++ b/lib/src/core/signal_client.dart @@ -1,4 +1,5 @@ import 'dart:async'; +import 'dart:collection'; import 'package:flutter_webrtc/flutter_webrtc.dart' as rtc; import 'package:http/http.dart' as http; @@ -26,6 +27,8 @@ class SignalClient extends Disposable with EventsEmittable { final WebSocketConnector _wsConnector; LiveKitWebSocket? _ws; + final _queue = Queue(); + @internal SignalClient(WebSocketConnector wsConnector) : _wsConnector = wsConnector { events.listen((event) { @@ -34,7 +37,7 @@ class SignalClient extends Disposable with EventsEmittable { onDispose(() async { await events.dispose(); - await _cleanUp(); + await cleanUp(); }); } @@ -59,7 +62,7 @@ class SignalClient extends Disposable with EventsEmittable { ? ConnectionState.reconnecting : ConnectionState.connecting); // Clean up existing socket - await _cleanUp(); + await cleanUp(); // Attempt to connect _ws = await _wsConnector( rtcUri, @@ -103,26 +106,41 @@ class SignalClient extends Disposable with EventsEmittable { } } - Future _cleanUp() async { + @internal + Future cleanUp() async { await _ws?.dispose(); _ws = null; + _queue.clear(); } @internal Future disconnect() async { logger.fine('SignalClient disconnect'); - await _cleanUp(); + await cleanUp(); } - void _sendRequest(lk_rtc.SignalRequest req) { - if (_ws == null || isDisposed) { - logger.warning( - '[$objectId] Could not send message, not connected or already disposed'); + void _sendRequest( + lk_rtc.SignalRequest req, { + bool enqueueIfReconnecting = true, + }) { + if (isDisposed) { + logger.warning('[$objectId] Could not send message, already disposed'); return; } - final buf = req.writeToBuffer(); - _ws?.send(buf); + if (_connectionState == ConnectionState.reconnecting && + req._canQueue() && + enqueueIfReconnecting) { + _queue.add(req); + return; + } + + if (_ws == null) { + logger.warning('[$objectId] Could not send message, socket is null'); + return; + } + + _ws?.send(req.writeToBuffer()); } void _updateConnectionState(ConnectionState newValue) { @@ -379,3 +397,34 @@ extension SignalClientRequests on SignalClient { ), )); } + +// private methods +extension on lk_rtc.SignalRequest { + // returns if this request can be queued + bool _canQueue() => ![ + // list of types that cannot be queued + lk_rtc.SignalRequest_Message.syncState, + lk_rtc.SignalRequest_Message.trickle, + lk_rtc.SignalRequest_Message.offer, + lk_rtc.SignalRequest_Message.answer, + lk_rtc.SignalRequest_Message.simulate + ].contains(whichMessage()); +} + +// internal methods + +extension SignalClientInternalMethods on SignalClient { + @internal + void sendQueuedRequests() { + // queue is empty + if (_queue.isEmpty) return; + // send requests + for (final request in _queue) { + _sendRequest(request, enqueueIfReconnecting: false); + } + _queue.clear(); + } + + @internal + void clearQueue() => _queue.clear(); +}