From 7acfb3faacf10b50ed12b0a8fa4e9e02b814e8dd Mon Sep 17 00:00:00 2001 From: CloudWebRTC Date: Tue, 18 Oct 2022 15:52:24 +0800 Subject: [PATCH] chore: Make timeouts configurable. (#191) --- lib/src/constants.dart | 26 +++++++++++++++++++++----- lib/src/core/engine.dart | 20 ++++++++++---------- lib/src/core/transport.dart | 12 +++++++----- lib/src/options.dart | 4 ++++ lib/src/participant/remote.dart | 3 +-- 5 files changed, 43 insertions(+), 22 deletions(-) diff --git a/lib/src/constants.dart b/lib/src/constants.dart index eaab5c4..bb28f32 100644 --- a/lib/src/constants.dart +++ b/lib/src/constants.dart @@ -1,7 +1,23 @@ class Timeouts { - static const connection = Duration(seconds: 10); - static const debounce = Duration(milliseconds: 100); - static const publish = Duration(seconds: 10); - static const peerConnection = Duration(seconds: 10); - static const iceRestart = Duration(seconds: 10); + final Duration connection; + final Duration debounce; + final Duration publish; + final Duration peerConnection; + final Duration iceRestart; + + const Timeouts({ + required this.connection, + required this.debounce, + required this.publish, + required this.peerConnection, + required this.iceRestart, + }); + + static const Timeouts defaultTimeouts = Timeouts( + connection: Duration(seconds: 10), + debounce: Duration(milliseconds: 100), + publish: Duration(seconds: 10), + peerConnection: Duration(seconds: 10), + iceRestart: Duration(seconds: 10), + ); } diff --git a/lib/src/core/engine.dart b/lib/src/core/engine.dart index 5e3f7d5..5c9e5cc 100644 --- a/lib/src/core/engine.dart +++ b/lib/src/core/engine.dart @@ -6,7 +6,6 @@ import 'package:collection/collection.dart'; import 'package:flutter_webrtc/flutter_webrtc.dart' as rtc; import 'package:meta/meta.dart'; -import '../constants.dart'; import '../events.dart'; import '../exceptions.dart'; import '../extensions.dart'; @@ -136,7 +135,7 @@ class Engine extends Disposable with EventsEmittable { // wait for join response await _signalListener.waitFor( - duration: Timeouts.connection, + duration: this.connectOptions.timeouts.connection, onTimeout: () => throw ConnectException(), ); @@ -145,7 +144,7 @@ class Engine extends Disposable with EventsEmittable { // wait until engine is connected await events.waitFor( filter: (event) => event.isPrimary && event.state.isConnected(), - duration: Timeouts.connection, + duration: this.connectOptions.timeouts.connection, onTimeout: () => throw ConnectException(), ); @@ -199,7 +198,7 @@ class Engine extends Disposable with EventsEmittable { // wait for response, or timeout final event = await _signalListener.waitFor( filter: (event) => event.cid == cid, - duration: Timeouts.publish, + duration: connectOptions.timeouts.publish, onTimeout: () => throw TrackPublishException(), ); @@ -242,7 +241,7 @@ class Engine extends Disposable with EventsEmittable { logger.fine('Waiting for publisher to ice-connect...'); await events.waitFor( filter: (event) => event.state.isConnected(), - duration: Timeouts.peerConnection, + duration: connectOptions.timeouts.peerConnection, ); } @@ -252,7 +251,7 @@ class Engine extends Disposable with EventsEmittable { logger.fine('Waiting for data channel ${reliability} to open...'); await events.waitFor( filter: (event) => event.type == reliability, - duration: Timeouts.connection, + duration: connectOptions.timeouts.connection, ); } } @@ -313,7 +312,7 @@ class Engine extends Disposable with EventsEmittable { await events.waitFor( filter: (event) => event.isPrimary && event.state.isConnected(), - duration: Timeouts.iceRestart, + duration: connectOptions.timeouts.iceRestart, onTimeout: () => throw ConnectException(), ); } @@ -361,9 +360,10 @@ class Engine extends Disposable with EventsEmittable { iceTransportPolicy: RTCIceTransportPolicy.relay); } - publisher = await Transport.create(_peerConnectionCreate, rtcConfiguration); - subscriber = - await Transport.create(_peerConnectionCreate, rtcConfiguration); + publisher = await Transport.create(_peerConnectionCreate, + rtcConfig: rtcConfiguration, connectOptions: connectOptions); + subscriber = await Transport.create(_peerConnectionCreate, + rtcConfig: rtcConfiguration, connectOptions: connectOptions); publisher?.pc.onIceCandidate = (rtc.RTCIceCandidate candidate) { logger.fine('publisher onIceCandidate'); diff --git a/lib/src/core/transport.dart b/lib/src/core/transport.dart index 2459bd7..1d74f29 100644 --- a/lib/src/core/transport.dart +++ b/lib/src/core/transport.dart @@ -1,8 +1,8 @@ import 'dart:async'; import 'package:flutter_webrtc/flutter_webrtc.dart' as rtc; +import 'package:livekit_client/src/options.dart'; -import '../constants.dart'; import '../extensions.dart'; import '../internal/types.dart'; import '../logger.dart'; @@ -23,9 +23,10 @@ class Transport extends Disposable { bool renegotiate = false; TransportOnOffer? onOffer; Function? _cancelDebounce; + ConnectOptions connectOptions; // private constructor - Transport._(this.pc) { + Transport._(this.pc, this.connectOptions) { // onDispose(() async { _cancelDebounce?.call(); @@ -59,17 +60,18 @@ class Transport extends Disposable { } static Future create(PeerConnectionCreate peerConnectionCreate, - [RTCConfiguration? rtcConfig]) async { + {RTCConfiguration? rtcConfig, + required ConnectOptions connectOptions}) async { rtcConfig ??= const RTCConfiguration(); logger.fine('[PCTransport] creating ${rtcConfig.toMap()}'); final pc = await peerConnectionCreate(rtcConfig.toMap()); - return Transport._(pc); + return Transport._(pc, connectOptions); } late final negotiate = Utils.createDebounceFunc( (void _) => createAndSendOffer(), cancelFunc: (f) => _cancelDebounce = f, - wait: Timeouts.debounce, + wait: connectOptions.timeouts.debounce, ); Future setRemoteDescription(rtc.RTCSessionDescription sd) async { diff --git a/lib/src/options.dart b/lib/src/options.dart index c2b61c0..5c6f87f 100644 --- a/lib/src/options.dart +++ b/lib/src/options.dart @@ -1,3 +1,4 @@ +import 'constants.dart'; import 'core/room.dart'; import 'publication/remote.dart'; import 'track/local/audio.dart'; @@ -38,10 +39,13 @@ class ConnectOptions { /// The protocol version to be used. Usually this doesn't need to be modified. final ProtocolVersion protocolVersion; + final Timeouts timeouts; + const ConnectOptions({ this.autoSubscribe = true, this.rtcConfiguration = const RTCConfiguration(), this.protocolVersion = ProtocolVersion.v8, + this.timeouts = Timeouts.defaultTimeouts, }); } diff --git a/lib/src/participant/remote.dart b/lib/src/participant/remote.dart index 7f8d24b..1791834 100644 --- a/lib/src/participant/remote.dart +++ b/lib/src/participant/remote.dart @@ -1,7 +1,6 @@ import 'package:flutter_webrtc/flutter_webrtc.dart' as rtc; import 'package:meta/meta.dart'; -import '../constants.dart'; import '../core/room.dart'; import '../events.dart'; import '../exceptions.dart'; @@ -88,7 +87,7 @@ class RemoteParticipant extends Participant { final event = await events.waitFor( filter: (event) => event.participant == this && event.publication.sid == trackSid, - duration: Timeouts.publish, + duration: room.connectOptions.timeouts.publish, onTimeout: () => throw TrackSubscriptionExceptionEvent( participant: this, sid: trackSid,