Session migration (#72)

* impl

* retry logic

* simulate

* pass checks

* adjust log level

* retry condition

* fix simulate scenario params

* format

* use `RTCPeerConnectionState` instead of `RTCIceConnectionState`

* adjust retry

* additional check for completer

* more logging

* fix id

* pub

* change package

* doc update

* retry tests

* format

* apply suggestions
This commit is contained in:
Hiroshi Horie
2022-01-25 00:53:18 +07:00
committed by GitHub
parent 3d37f4491d
commit 224047c31f
25 changed files with 575 additions and 275 deletions
+142 -114
View File
@@ -19,6 +19,7 @@ import '../proto/livekit_models.pb.dart' as lk_models;
import '../proto/livekit_rtc.pb.dart' as lk_rtc;
import '../support/disposable.dart';
import '../types.dart';
import '../utils.dart';
import 'room.dart';
import 'signal_client.dart';
import 'transport.dart';
@@ -26,7 +27,6 @@ import 'transport.dart';
class Engine extends Disposable with EventsEmittable<EngineEvent> {
static const _lossyDCLabel = '_lossy';
static const _reliableDCLabel = '_reliable';
static const _maxReconnectAttempts = 5;
// Reference to the Room
final Room room;
@@ -48,8 +48,6 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
rtc.RTCDataChannel? _reliableDCSub;
rtc.RTCDataChannel? _lossyDCSub;
bool _iceConnected = false;
ConnectionState _connectionState = ConnectionState.disconnected;
/// Connection state of the [Room].
@@ -68,9 +66,6 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
// server-provided ice servers
List<lk_rtc.ICEServer> _serverProvidedIceServers = [];
// internal
int _reconnectAttempts = 0;
late final _signalListener = signalClient.createListener(synchronized: true);
final delays = CancelableDelayManager();
@@ -94,34 +89,51 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
});
}
Future<lk_rtc.JoinResponse> connect(
Future<void> connect(
String url,
String token,
) async {
this.url = url;
this.token = token;
// connect to rtc server
await signalClient.connect(
url,
token,
connectOptions: room.connectOptions,
);
_updateConnectionState(ConnectionState.connecting);
// wait for join response
final event = await _signalListener.waitFor<SignalConnectedEvent>(
duration: Timeouts.connection,
onTimeout: () => throw ConnectException(),
);
try {
// wait for socket to connect rtc server
await signalClient.connect(
url,
token,
connectOptions: room.connectOptions,
);
return event.response;
// wait for join response
await _signalListener.waitFor<SignalJoinResponseEvent>(
duration: Timeouts.connection,
onTimeout: () => throw ConnectException(),
);
logger.fine('Waiting for engine to connect...');
// wait until engine is connected
await events.waitFor<EnginePeerStateUpdatedEvent>(
filter: (event) => event.isPrimary && event.state.isConnected(),
duration: Timeouts.connection,
onTimeout: () => throw ConnectException(),
);
_updateConnectionState(ConnectionState.connected);
} catch (error) {
logger.fine('Connect Error $error');
_updateConnectionState(ConnectionState.disconnected);
rethrow;
}
}
/// Close connection between the server.
Future<void> close() async {
logger.fine('[$objectId] close()');
logger.fine('${runtimeType}.close()');
if (_connectionState == ConnectionState.disconnected) {
logger.warning('[$objectId]: close() already disconnected');
logger.warning('${runtimeType}.close() already disconnected');
}
// _statsTimer.cancel();
// cancel all ongoing delays
@@ -134,9 +146,9 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
await subscriber?.dispose();
subscriber = null;
await signalClient.close();
await signalClient.disconnect();
_connectionState = ConnectionState.disconnected;
_updateConnectionState(ConnectionState.disconnected);
// notifyListeners();
}
@@ -148,7 +160,7 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
required lk_models.TrackSource source,
VideoDimensions? dimensions,
bool? dtx,
List<lk_models.VideoLayer>? videoLayers,
Iterable<lk_models.VideoLayer>? videoLayers,
}) async {
// TODO: Check if cid already published
@@ -205,19 +217,19 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
if (_subscriberPrimary) {
// make sure publisher transport is connected
if (publisher?.pc.iceConnectionState?.isConnected() != true) {
if (publisher?.pc.connectionState?.isConnected() != true) {
logger.fine('Publisher is not connected...');
// start negotiation
if (publisher?.pc.iceConnectionState !=
rtc.RTCIceConnectionState.RTCIceConnectionStateChecking) {
if (publisher?.pc.connectionState !=
rtc.RTCPeerConnectionState.RTCPeerConnectionStateConnecting) {
await negotiate();
}
logger.fine('Waiting for publisher to ice-connect...');
await events.waitFor<EnginePublisherIceStateUpdatedEvent>(
filter: (event) => event.iceState.isConnected(),
duration: Timeouts.iceConnection,
await events.waitFor<EnginePublisherPeerStateUpdatedEvent>(
filter: (event) => event.state.isConnected(),
duration: Timeouts.peerConnection,
);
}
@@ -251,25 +263,17 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
return;
}
final url = this.url;
final token = this.token;
if (url == null || token == null) {
throw ConnectException('could not reconnect without url and token');
}
_connectionState = ConnectionState.reconnecting;
if (_reconnectAttempts == 0) {
events.emit(const EngineReconnectingEvent());
}
_reconnectAttempts++;
try {
await signalClient.reconnect(
url,
token,
Future<void> sequence() async {
//
await signalClient.connect(
url!,
token!,
connectOptions: room.connectOptions,
reconnect: true,
);
if (publisher == null || subscriber == null) {
@@ -285,29 +289,37 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
));
}
final iceConnected =
primary?.pc.iceConnectionState?.isConnected() ?? false;
final iceConnected = primary?.pc.connectionState?.isConnected() ?? false;
logger.fine('Reconnect: iceConnected: $iceConnected');
if (!iceConnected) {
logger.fine('Reconnect: Waiting for primary to connect...');
await events.waitFor<EngineIceStateUpdatedEvent>(
filter: (event) => event.isPrimary && event.iceState.isConnected(),
await events.waitFor<EnginePeerStateUpdatedEvent>(
filter: (event) => event.isPrimary && event.state.isConnected(),
duration: Timeouts.iceRestart,
onTimeout: () => throw ConnectException(),
);
}
}
logger.fine('Reconnect: success');
_connectionState = ConnectionState.connected;
events.emit(const EngineReconnectedEvent());
_reconnectAttempts = 0;
try {
_updateConnectionState(ConnectionState.reconnecting);
await Utils.retry<void>(
(tries, errors) {
logger.fine('Retrying connect sequence remaining ${tries} tries...');
return sequence();
},
retryCondition: (_, __) =>
_connectionState == ConnectionState.reconnecting,
tries: 3,
delay: const Duration(seconds: 3),
);
_updateConnectionState(ConnectionState.connected);
} catch (error) {
logger.fine('Reconnect: error ${error}');
// Pass up all exceptions
_connectionState = ConnectionState.disconnected;
rethrow;
//
_updateConnectionState(ConnectionState.disconnected);
}
}
@@ -353,39 +365,26 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
subscriber?.pc.onDataChannel = _onDataChannel;
}
subscriber?.pc.onIceConnectionState =
(state) => events.emit(EngineSubscriberIceStateUpdatedEvent(
iceState: state,
subscriber?.pc.onConnectionState =
(state) => events.emit(EngineSubscriberPeerStateUpdatedEvent(
state: state,
isPrimary: _subscriberPrimary,
));
publisher?.pc.onIceConnectionState =
(state) => events.emit(EnginePublisherIceStateUpdatedEvent(
publisher?.pc.onConnectionState =
(state) => events.emit(EnginePublisherPeerStateUpdatedEvent(
state: state,
isPrimary: !_subscriberPrimary,
));
events.on<EngineIceStateUpdatedEvent>((event) {
events.on<EnginePeerStateUpdatedEvent>((event) {
// only listen to primary ice events
if (!event.isPrimary) return;
if (event.iceState ==
rtc.RTCIceConnectionState.RTCIceConnectionStateConnected) {
if (!_iceConnected) {
_iceConnected = true;
if (_connectionState == ConnectionState.reconnecting) {
events.emit(const EngineReconnectedEvent());
} else {
events.emit(const EngineConnectedEvent());
}
}
} else if (event.iceState ==
rtc.RTCIceConnectionState.RTCIceConnectionStateFailed) {
if (event.state ==
rtc.RTCPeerConnectionState.RTCPeerConnectionStateFailed) {
// trigger reconnect sequence
if (_iceConnected) {
_iceConnected = false;
_onDisconnected('peerconnection');
}
_onDisconnected(DisconnectReason.peerConnection);
}
});
@@ -502,8 +501,9 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
final dp = lk_models.DataPacket.fromBuffer(message.binary);
if (dp.whichValue() == lk_models.DataPacket_Value.speaker) {
// Speaker packet
events
.emit(EngineActiveSpeakersUpdateEvent(speakers: dp.speaker.speakers));
events.emit(EngineActiveSpeakersUpdateEvent(
speakers: dp.speaker.speakers,
));
} else if (dp.whichValue() == lk_models.DataPacket_Value.user) {
// User packet
events.emit(EngineDataPacketReceivedEvent(
@@ -513,44 +513,63 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
}
}
Future<void> _onDisconnected(String reason) async {
logger.info('onDisconnected reason: $reason');
Future<void> _onDisconnected(DisconnectReason reason) async {
logger
.info('onDisconnected state:${_connectionState} reason:${reason.name}');
if (_connectionState == ConnectionState.disconnected) {
logger.fine('[$objectId] Already disconnected $reason');
logger.fine('[$objectId] Already disconnected... $reason');
return;
}
if (_connectionState == ConnectionState.reconnecting) {
logger.fine('[$objectId] Already reconnecting...');
return;
}
logger.fine('[$objectId] Disconnected $reason');
logger.fine('[$runtimeType] Should attempt reconnect sequence...');
if (_reconnectAttempts >= _maxReconnectAttempts) {
logger.info('[$objectId] Could not connect '
'after ${_reconnectAttempts} attempts, giving up');
await close();
events.emit(const EngineDisconnectedEvent());
return;
}
await reconnect();
}
final delay =
Duration(milliseconds: (_reconnectAttempts * _reconnectAttempts) * 300);
void _updateConnectionState(ConnectionState newValue) {
if (_connectionState == newValue) return;
// if this instance is disposed, we probably don't want to continue any more
// so the whole block will be canceled from being executed
await delays.waitFor(delay, ifNotCancelled: () async {
try {
await reconnect();
_reconnectAttempts = 0;
} catch (_) {
// doesn't need to be awaited
// ignore: unawaited_futures
_onDisconnected(reason);
logger.fine('Engine ConnectionState '
'${_connectionState.name} -> ${newValue.name}');
bool didReconnect = _connectionState == ConnectionState.reconnecting &&
newValue == ConnectionState.connected;
// update internal value
_connectionState = newValue;
// emit event
if (_connectionState == ConnectionState.connected) {
if (didReconnect) {
events.emit(const EngineReconnectedEvent());
} else {
events.emit(const EngineConnectedEvent());
}
});
} else if (_connectionState == ConnectionState.reconnecting) {
events.emit(const EngineReconnectingEvent());
} else if (_connectionState == ConnectionState.disconnected) {
events.emit(const EngineDisconnectedEvent());
}
}
@internal
void sendSyncState({
required lk_rtc.UpdateSubscription subscription,
required Iterable<lk_rtc.TrackPublishedResponse>? publishTracks,
}) async {
final answer = (await subscriber?.pc.getLocalDescription())?.toPBType();
signalClient.sendSyncState(
answer: answer,
subscription: subscription,
publishTracks: publishTracks,
);
}
void _setUpListeners() => _signalListener
..on<SignalConnectedEvent>((event) async {
..on<SignalJoinResponseEvent>((event) async {
// create peer connections
_connectionState = ConnectionState.connected;
_subscriberPrimary = event.response.subscriberPrimary;
_serverProvidedIceServers = event.response.iceServers;
@@ -564,9 +583,16 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
// for subscriberPrimary, we negotiate when necessary (lazy)
await negotiate();
}
// Relay to Room
events.emit(event);
})
..on<SignalCloseEvent>((_) async {
await _onDisconnected('signal');
..on<SignalConnectionStateUpdatedEvent>((event) async {
if (event.connectionState == ConnectionState.disconnected) {
await _onDisconnected(DisconnectReason.signal);
}
// Relay to Room
events.emit(event);
})
..on<SignalOfferEvent>((event) async {
if (subscriber == null) {
@@ -611,21 +637,23 @@ class Engine extends Disposable with EventsEmittable<EngineEvent> {
await publisher!.addIceCandidate(event.candidate);
}
})
// relay
// relay to Room
..on<SignalParticipantUpdateEvent>((event) => events.emit(event))
// relay
// relay to Room
..on<SignalSpeakersChangedEvent>((event) => events.emit(event))
// relay
// relay to Room
..on<SignalConnectionQualityUpdateEvent>((event) => events.emit(event))
// relay
// relay to Room
..on<SignalStreamStateUpdatedEvent>((event) => events.emit(event))
// relay to Room
..on<SignalSubscribedQualityUpdatedEvent>((event) => events.emit(event))
// relay to Room
..on<SignalSubscriptionPermissionUpdateEvent>((event) => events.emit(event))
..on<SignalLeaveEvent>((event) async {
if (connectionState == ConnectionState.reconnecting) {
logger.warning('Received leave signal while engine is reconnecting.');
if (_connectionState == ConnectionState.reconnecting) {
logger.warning(
'[Signal] Received Leave while engine is reconnecting, ignoring...');
return;
}
await close();
events.emit(const EngineDisconnectedEvent());