Enqueue Signal requests while reconnecting (#81)

* implement queue to events emitter

* check `isDisposed` instead

* update message

* implement queue to events emitter

* update lock

* impl
This commit is contained in:
Hiroshi Horie
2022-03-02 08:14:23 +09:00
committed by GitHub
parent da8fe67ec6
commit 811ded7aa6
2 changed files with 62 additions and 10 deletions
+3
View File
@@ -545,12 +545,15 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
if (_connectionState == ConnectionState.connected) { if (_connectionState == ConnectionState.connected) {
if (didReconnect) { if (didReconnect) {
events.emit(const EngineReconnectedEvent()); events.emit(const EngineReconnectedEvent());
// send queued requests if engine re-connected
signalClient.sendQueuedRequests();
} else { } else {
events.emit(const EngineConnectedEvent()); events.emit(const EngineConnectedEvent());
} }
} else if (_connectionState == ConnectionState.reconnecting) { } else if (_connectionState == ConnectionState.reconnecting) {
events.emit(const EngineReconnectingEvent()); events.emit(const EngineReconnectingEvent());
} else if (_connectionState == ConnectionState.disconnected) { } else if (_connectionState == ConnectionState.disconnected) {
signalClient.cleanUp();
events.emit(const EngineDisconnectedEvent()); events.emit(const EngineDisconnectedEvent());
} }
} }
+59 -10
View File
@@ -1,4 +1,5 @@
import 'dart:async'; import 'dart:async';
import 'dart:collection';
import 'package:flutter_webrtc/flutter_webrtc.dart' as rtc; import 'package:flutter_webrtc/flutter_webrtc.dart' as rtc;
import 'package:http/http.dart' as http; import 'package:http/http.dart' as http;
@@ -26,6 +27,8 @@ class SignalClient extends Disposable with EventsEmittable<SignalEvent> {
final WebSocketConnector _wsConnector; final WebSocketConnector _wsConnector;
LiveKitWebSocket? _ws; LiveKitWebSocket? _ws;
final _queue = Queue<lk_rtc.SignalRequest>();
@internal @internal
SignalClient(WebSocketConnector wsConnector) : _wsConnector = wsConnector { SignalClient(WebSocketConnector wsConnector) : _wsConnector = wsConnector {
events.listen((event) { events.listen((event) {
@@ -34,7 +37,7 @@ class SignalClient extends Disposable with EventsEmittable<SignalEvent> {
onDispose(() async { onDispose(() async {
await events.dispose(); await events.dispose();
await _cleanUp(); await cleanUp();
}); });
} }
@@ -59,7 +62,7 @@ class SignalClient extends Disposable with EventsEmittable<SignalEvent> {
? ConnectionState.reconnecting ? ConnectionState.reconnecting
: ConnectionState.connecting); : ConnectionState.connecting);
// Clean up existing socket // Clean up existing socket
await _cleanUp(); await cleanUp();
// Attempt to connect // Attempt to connect
_ws = await _wsConnector( _ws = await _wsConnector(
rtcUri, rtcUri,
@@ -103,26 +106,41 @@ class SignalClient extends Disposable with EventsEmittable<SignalEvent> {
} }
} }
Future<void> _cleanUp() async { @internal
Future<void> cleanUp() async {
await _ws?.dispose(); await _ws?.dispose();
_ws = null; _ws = null;
_queue.clear();
} }
@internal @internal
Future<void> disconnect() async { Future<void> disconnect() async {
logger.fine('SignalClient disconnect'); logger.fine('SignalClient disconnect');
await _cleanUp(); await cleanUp();
} }
void _sendRequest(lk_rtc.SignalRequest req) { void _sendRequest(
if (_ws == null || isDisposed) { lk_rtc.SignalRequest req, {
logger.warning( bool enqueueIfReconnecting = true,
'[$objectId] Could not send message, not connected or already disposed'); }) {
if (isDisposed) {
logger.warning('[$objectId] Could not send message, already disposed');
return; return;
} }
final buf = req.writeToBuffer(); if (_connectionState == ConnectionState.reconnecting &&
_ws?.send(buf); 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) { 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();
}