[frdp] Add configurable timeouts for VM. (#22893) This adds configurable timeouts for the Dart VM. Due to some testing machines running things quite slowly, this is becoming more necessary.
diff --git a/packages/fuchsia_remote_debug_protocol/lib/src/dart/dart_vm.dart b/packages/fuchsia_remote_debug_protocol/lib/src/dart/dart_vm.dart index a2a6063..56b2abe 100644 --- a/packages/fuchsia_remote_debug_protocol/lib/src/dart/dart_vm.dart +++ b/packages/fuchsia_remote_debug_protocol/lib/src/dart/dart_vm.dart
@@ -18,9 +18,12 @@ final Logger _log = Logger('DartVm'); -/// Signature of an asynchronous function for astablishing a JSON RPC-2 +/// Signature of an asynchronous function for establishing a JSON RPC-2 /// connection to a [Uri]. -typedef RpcPeerConnectionFunction = Future<json_rpc.Peer> Function(Uri uri); +typedef RpcPeerConnectionFunction = Future<json_rpc.Peer> Function( + Uri uri, { + Duration timeout, +}); /// [DartVm] uses this function to connect to the Dart VM on Fuchsia. /// @@ -30,16 +33,18 @@ /// Attempts to connect to a Dart VM service. /// -/// Gives up after `_kConnectTimeout` has elapsed. -Future<json_rpc.Peer> _waitAndConnect(Uri uri) async { +/// Gives up after `timeout` has elapsed. +Future<json_rpc.Peer> _waitAndConnect( + Uri uri, { + Duration timeout = _kConnectTimeout, +}) async { final Stopwatch timer = Stopwatch()..start(); Future<json_rpc.Peer> attemptConnection(Uri uri) async { WebSocket socket; json_rpc.Peer peer; try { - socket = - await WebSocket.connect(uri.toString()).timeout(_kConnectTimeout); + socket = await WebSocket.connect(uri.toString()).timeout(timeout); peer = json_rpc.Peer(IOWebSocketChannel(socket).cast())..listen(); return peer; } on HttpException catch (e) { @@ -53,7 +58,7 @@ // Other unknown errors will be handled with reconnects. await peer?.close(); await socket?.close(); - if (timer.elapsed < _kConnectTimeout) { + if (timer.elapsed < timeout) { _log.info('Attempting to reconnect'); await Future<void>.delayed(_kReconnectAttemptInterval); return attemptConnection(uri); @@ -105,11 +110,15 @@ /// Attempts to connect to the given [Uri]. /// /// Throws an error if unable to connect. - static Future<DartVm> connect(Uri uri) async { + static Future<DartVm> connect( + Uri uri, { + Duration timeout = _kConnectTimeout, + }) async { if (uri.scheme == 'http') { uri = uri.replace(scheme: 'ws', path: '/ws'); } - final json_rpc.Peer peer = await fuchsiaVmServiceConnectionFunction(uri); + final json_rpc.Peer peer = + await fuchsiaVmServiceConnectionFunction(uri, timeout: timeout); if (peer == null) { return null; }
diff --git a/packages/fuchsia_remote_debug_protocol/lib/src/fuchsia_remote_connection.dart b/packages/fuchsia_remote_debug_protocol/lib/src/fuchsia_remote_connection.dart index fcbeb05..482e313 100644 --- a/packages/fuchsia_remote_debug_protocol/lib/src/fuchsia_remote_connection.dart +++ b/packages/fuchsia_remote_debug_protocol/lib/src/fuchsia_remote_connection.dart
@@ -20,6 +20,8 @@ const Duration _kIsolateFindTimeout = Duration(minutes: 1); +const Duration _kDartVmConnectionTimeout = Duration(seconds: 9); + const Duration _kVmPollInterval = Duration(milliseconds: 1500); final Logger _log = Logger('FuchsiaRemoteConnection'); @@ -246,6 +248,7 @@ Future<List<IsolateRef>> _waitForMainIsolatesByPattern([ Pattern pattern, Duration timeout = _kIsolateFindTimeout, + Duration vmConnectionTimeout = _kDartVmConnectionTimeout, ]) async { final Completer<List<IsolateRef>> completer = Completer<List<IsolateRef>>(); _onDartVmEvent.listen( @@ -253,7 +256,8 @@ if (event.eventType == DartVmEventType.started) { _log.fine('New VM found on port: ${event.servicePort}. Searching ' 'for Isolate: $pattern'); - final DartVm vmService = await _getDartVm(event.uri.port); + final DartVm vmService = await _getDartVm(event.uri.port, + timeout: _kDartVmConnectionTimeout); // If the VM service is null, set the result to the empty list. final List<IsolateRef> result = await vmService ?.getMainIsolatesByPattern(pattern, timeout: timeout) ?? @@ -284,29 +288,33 @@ /// either `timeout` is reached, or a Dart VM starts up with a name that /// matches `pattern`. Future<List<IsolateRef>> getMainIsolatesByPattern( - Pattern pattern, [ + Pattern pattern, { Duration timeout = _kIsolateFindTimeout, - ]) async { + Duration vmConnectionTimeout = _kDartVmConnectionTimeout, + }) async { // If for some reason there are no Dart VM's that are alive, wait for one to // start with the Isolate in question. if (_dartVmPortMap.isEmpty) { _log.fine('No live Dart VMs found. Awaiting new VM startup'); - return _waitForMainIsolatesByPattern(pattern, timeout); + return _waitForMainIsolatesByPattern( + pattern, timeout, vmConnectionTimeout); } // Accumulate a list of eventual IsolateRef lists so that they can be loaded // simultaneously via Future.wait. final List<Future<List<IsolateRef>>> isolates = <Future<List<IsolateRef>>>[]; for (PortForwarder fp in _dartVmPortMap.values) { - final DartVm vmService = await _getDartVm(fp.port).timeout(timeout); + final DartVm vmService = + await _getDartVm(fp.port, timeout: vmConnectionTimeout); if (vmService == null) { continue; } isolates.add(vmService.getMainIsolatesByPattern(pattern)); } - final List<IsolateRef> result = await Future.wait<List<IsolateRef>>(isolates) - .timeout(timeout) - .then<List<IsolateRef>>((List<List<IsolateRef>> listOfLists) { + final List<IsolateRef> result = + await Future.wait<List<IsolateRef>>(isolates) + .timeout(timeout) + .then<List<IsolateRef>>((List<List<IsolateRef>> listOfLists) { final List<List<IsolateRef>> mutableListOfLists = List<List<IsolateRef>>.from(listOfLists) ..retainWhere((List<IsolateRef> list) => list.isNotEmpty); @@ -330,7 +338,8 @@ // TODO(awdavies): Set this up to handle multiple Isolates per Dart VM. if (result.isEmpty) { _log.fine('No instance of the Isolate found. Awaiting new VM startup'); - return _waitForMainIsolatesByPattern(pattern, timeout); + return _waitForMainIsolatesByPattern( + pattern, timeout, vmConnectionTimeout); } return result; } @@ -407,13 +416,17 @@ /// /// Returns null if either there is an [HttpException] or a /// [TimeoutException], else a [DartVm] instance. - Future<DartVm> _getDartVm(int port) async { + Future<DartVm> _getDartVm( + int port, { + Duration timeout = _kDartVmConnectionTimeout, + }) async { if (!_dartVmCache.containsKey(port)) { // When raising an HttpException this means that there is no instance of // the Dart VM to communicate with. The TimeoutException is raised when // the Dart VM instance is shut down in the middle of communicating. try { - final DartVm dartVm = await DartVm.connect(_getDartVmUri(port)); + final DartVm dartVm = + await DartVm.connect(_getDartVmUri(port), timeout: timeout); _dartVmCache[port] = dartVm; } on HttpException { _log.warning('HTTP Exception encountered connecting to new VM'); @@ -460,7 +473,7 @@ (DartVm vmService) async { final Map<String, dynamic> res = await vmService.invokeRpc('getVersion'); - _log.fine('DartVM version check result: $res'); + _log.fine('DartVM(${vmService.uri}) version check result: $res'); return res; }, queueEvents, @@ -475,7 +488,8 @@ await stop(); final List<int> servicePorts = await getDeviceServicePorts(); final List<PortForwarder> forwardedVmServicePorts = - await Future.wait<PortForwarder>(servicePorts.map<Future<PortForwarder>>((int deviceServicePort) { + await Future.wait<PortForwarder>( + servicePorts.map<Future<PortForwarder>>((int deviceServicePort) { return fuchsiaPortForwardingFunction( _sshCommandRunner.address, deviceServicePort,
diff --git a/packages/fuchsia_remote_debug_protocol/test/fuchsia_remote_connection_test.dart b/packages/fuchsia_remote_debug_protocol/test/fuchsia_remote_connection_test.dart index b282f17..dc62fdf 100644 --- a/packages/fuchsia_remote_debug_protocol/test/fuchsia_remote_connection_test.dart +++ b/packages/fuchsia_remote_debug_protocol/test/fuchsia_remote_connection_test.dart
@@ -22,8 +22,7 @@ mockRunner = MockSshCommandRunner(); // Adds some extra junk to make sure the strings will be cleaned up. when(mockRunner.run(any)).thenAnswer((_) => - Future<List<String>>.value( - <String>['123\n\n\n', '456 ', '789'])); + Future<List<String>>.value(<String>['123\n\n\n', '456 ', '789'])); const String address = 'fe80::8eae:4cff:fef4:9247'; const String interface = 'eno1'; when(mockRunner.address).thenReturn(address); @@ -86,7 +85,10 @@ mockPeerConnections = <MockPeer>[]; uriConnections = <Uri>[]; - Future<json_rpc.Peer> mockVmConnectionFunction(Uri uri) { + Future<json_rpc.Peer> mockVmConnectionFunction( + Uri uri, { + Duration timeout, + }) { return Future<json_rpc.Peer>(() async { final MockPeer mp = MockPeer(); mockPeerConnections.add(mp);
diff --git a/packages/fuchsia_remote_debug_protocol/test/src/dart/dart_vm_test.dart b/packages/fuchsia_remote_debug_protocol/test/src/dart/dart_vm_test.dart index af000ec..7c5a07c 100644 --- a/packages/fuchsia_remote_debug_protocol/test/src/dart/dart_vm_test.dart +++ b/packages/fuchsia_remote_debug_protocol/test/src/dart/dart_vm_test.dart
@@ -17,7 +17,10 @@ }); test('null connector', () async { - Future<json_rpc.Peer> mockServiceFunction(Uri uri) { + Future<json_rpc.Peer> mockServiceFunction( + Uri uri, { + Duration timeout, + }) { return Future<json_rpc.Peer>(() => null); } @@ -28,7 +31,10 @@ test('disconnect closes peer', () async { final MockPeer peer = MockPeer(); - Future<json_rpc.Peer> mockServiceFunction(Uri uri) { + Future<json_rpc.Peer> mockServiceFunction( + Uri uri, { + Duration timeout, + }) { return Future<json_rpc.Peer>(() => peer); } @@ -84,7 +90,10 @@ ], }; - Future<json_rpc.Peer> mockVmConnectionFunction(Uri uri) { + Future<json_rpc.Peer> mockVmConnectionFunction( + Uri uri, { + Duration timeout, + }) { when(mockPeer.sendRequest(any, any)).thenAnswer((_) => Future<Map<String, dynamic>>(() => flutterViewCannedResponses)); return Future<json_rpc.Peer>(() => mockPeer); @@ -139,7 +148,10 @@ ], }; - Future<json_rpc.Peer> mockVmConnectionFunction(Uri uri) { + Future<json_rpc.Peer> mockVmConnectionFunction( + Uri uri, { + Duration timeout, + }) { when(mockPeer.sendRequest(any, any)).thenAnswer((_) => Future<Map<String, dynamic>>(() => flutterViewCannedResponses)); return Future<json_rpc.Peer>(() => mockPeer); @@ -186,7 +198,10 @@ ] }; - Future<json_rpc.Peer> mockVmConnectionFunction(Uri uri) { + Future<json_rpc.Peer> mockVmConnectionFunction( + Uri uri, { + Duration timeout, + }) { when(mockPeer.sendRequest(any, any)).thenAnswer((_) => Future<Map<String, dynamic>>( () => flutterViewCannedResponseMissingId)); @@ -239,7 +254,10 @@ ], }; - Future<json_rpc.Peer> mockVmConnectionFunction(Uri uri) { + Future<json_rpc.Peer> mockVmConnectionFunction( + Uri uri, { + Duration timeout, + }) { when(mockPeer.sendRequest(any, any)).thenAnswer( (_) => Future<Map<String, dynamic>>(() => vmCannedResponse)); return Future<json_rpc.Peer>(() => mockPeer); @@ -278,7 +296,10 @@ ], }; - Future<json_rpc.Peer> mockVmConnectionFunction(Uri uri) { + Future<json_rpc.Peer> mockVmConnectionFunction( + Uri uri, { + Duration timeout, + }) { when(mockPeer.sendRequest(any, any)).thenAnswer((_) => Future<Map<String, dynamic>>( () => flutterViewCannedResponseMissingIsolateName)); @@ -311,7 +332,10 @@ test('verify timeout fires', () async { const Duration timeoutTime = Duration(milliseconds: 100); - Future<json_rpc.Peer> mockVmConnectionFunction(Uri uri) { + Future<json_rpc.Peer> mockVmConnectionFunction( + Uri uri, { + Duration timeout, + }) { // Return a command that will never complete. when(mockPeer.sendRequest(any, any)) .thenAnswer((_) => Completer<Map<String, dynamic>>().future);