import 'dart:async';
import 'dart:convert';
import 'dart:developer';
import 'dart:io';
import 'dart:isolate';
import 'package:logging/logging.dart';
import 'package:stack_trace/stack_trace.dart';
import 'running_processes.dart';
import 'utils.dart';
typedef TaskFunction = Future<TaskResult> Function();
bool _isTaskRegistered = false;
Future<TaskResult> task(TaskFunction task) {
if (_isTaskRegistered)
throw StateError('A task is already registered');
_isTaskRegistered = true;
Logger.root.level = Level.ALL;
Logger.root.onRecord.listen((LogRecord rec) {
print('${rec.level.name}: ${rec.time}: ${rec.message}');
});
final _TaskRunner runner = _TaskRunner(task);
runner.keepVmAliveUntilTaskRunRequested();
return runner.whenDone;
}
class _TaskRunner {
_TaskRunner(this.task) {
registerExtension('ext.cocoonRunTask',
(String method, Map<String, String> parameters) async {
final Duration taskTimeout = parameters.containsKey('timeoutInMinutes')
? Duration(minutes: int.parse(parameters['timeoutInMinutes']))
: null;
final TaskResult result = await run(taskTimeout);
return ServiceExtensionResponse.result(json.encode(result.toJson()));
});
registerExtension('ext.cocoonRunnerReady',
(String method, Map<String, String> parameters) async {
return ServiceExtensionResponse.result('"ready"');
});
}
final TaskFunction task;
RawReceivePort _keepAlivePort;
Timer _startTaskTimeout;
bool _taskStarted = false;
final Completer<TaskResult> _completer = Completer<TaskResult>();
static final Logger logger = Logger('TaskRunner');
Future<TaskResult> get whenDone => _completer.future;
Future<TaskResult> run(Duration taskTimeout) async {
try {
_taskStarted = true;
print('Running task.');
final String exe = Platform.isWindows ? '.exe' : '';
section('Checking running Dart$exe processes');
final Set<RunningProcessInfo> beforeRunningDartInstances = await getRunningProcesses(
processName: 'dart$exe',
).toSet();
beforeRunningDartInstances.forEach(print);
Future<TaskResult> futureResult = _performTask();
if (taskTimeout != null)
futureResult = futureResult.timeout(taskTimeout);
TaskResult result = await futureResult;
section('Checking running Dart$exe processes after task...');
final List<RunningProcessInfo> afterRunningDartInstances = await getRunningProcesses(
processName: 'dart$exe',
).toList();
for (final RunningProcessInfo info in afterRunningDartInstances) {
if (!beforeRunningDartInstances.contains(info)) {
print('$info was leaked by this test.');
if (result is TaskResultCheckProcesses) {
result = TaskResult.failure('This test leaked dart processes');
}
final bool killed = await killProcess(info.pid);
if (!killed) {
print('Failed to kill process ${info.pid}.');
} else {
print('Killed process id ${info.pid}.');
}
}
}
_completer.complete(result);
return result;
} on TimeoutException catch (_) {
print('Task timed out in framework.dart after $taskTimeout.');
return TaskResult.failure('Task timed out after $taskTimeout');
} finally {
print('Cleaning up after task...');
await forceQuitRunningProcesses();
_closeKeepAlivePort();
}
}
void keepVmAliveUntilTaskRunRequested() {
if (_taskStarted)
throw StateError('Task already started.');
_keepAlivePort = RawReceivePort();
const Duration taskStartTimeout = Duration(seconds: 60);
_startTaskTimeout = Timer(taskStartTimeout, () {
if (!_taskStarted) {
logger.severe('Task did not start in $taskStartTimeout.');
_closeKeepAlivePort();
exitCode = 1;
}
});
}
void _closeKeepAlivePort() {
_startTaskTimeout?.cancel();
_keepAlivePort?.close();
}
Future<TaskResult> _performTask() {
final Completer<TaskResult> completer = Completer<TaskResult>();
Chain.capture(() async {
completer.complete(await task());
}, onError: (dynamic taskError, Chain taskErrorStack) {
final String message = 'Task failed: $taskError';
stderr
..writeln(message)
..writeln('\nStack trace:')
..writeln(taskErrorStack.terse);
if (!completer.isCompleted)
completer.complete(TaskResult.failure(message));
});
return completer.future;
}
}
class TaskResult {
TaskResult.success(this.data, {this.benchmarkScoreKeys = const <String>[]})
: succeeded = true,
message = 'success' {
const JsonEncoder prettyJson = JsonEncoder.withIndent(' ');
if (benchmarkScoreKeys != null) {
for (String key in benchmarkScoreKeys) {
if (!data.containsKey(key)) {
throw 'Invalid Golem score key "$key". It does not exist in task '
'result data ${prettyJson.convert(data)}';
} else if (data[key] is! num) {
throw 'Invalid Golem score for key "$key". It is expected to be a num '
'but was ${data[key].runtimeType}: ${prettyJson.convert(data[key])}';
}
}
}
}
factory TaskResult.successFromFile(File file,
{List<String> benchmarkScoreKeys}) {
return TaskResult.success(json.decode(file.readAsStringSync()),
benchmarkScoreKeys: benchmarkScoreKeys);
}
TaskResult.failure(this.message)
: succeeded = false,
data = null,
benchmarkScoreKeys = const <String>[];
final bool succeeded;
final Map<String, dynamic> data;
final List<String> benchmarkScoreKeys;
bool get failed => !succeeded;
final String message;
Map<String, dynamic> toJson() {
final Map<String, dynamic> json = <String, dynamic>{
'success': succeeded,
};
if (succeeded) {
json['data'] = data;
json['benchmarkScoreKeys'] = benchmarkScoreKeys;
} else {
json['reason'] = message;
}
return json;
}
}
class TaskResultCheckProcesses extends TaskResult {
TaskResultCheckProcesses() : super.success(null);
}