spawn static method
Future<NativeTaskIsolate>
spawn({
- required NativeTaskWorkerFactory factory,
- required Object? initialMessage,
- String? debugName,
Starts a worker with factory and its sendable initialMessage.
Implementation
static Future<NativeTaskIsolate> spawn({
required NativeTaskWorkerFactory factory,
required Object? initialMessage,
String? debugName,
}) async {
final ReceivePort responses = ReceivePort();
final ReceivePort errors = ReceivePort();
final ReceivePort exits = ReceivePort();
final Completer<_NativeTaskReady> ready = Completer<_NativeTaskReady>();
NativeTaskIsolate? worker;
Object? earlyTermination;
late final StreamSubscription<Object?> responseSubscription;
responseSubscription = responses.listen((Object? message) {
if (message case final _NativeTaskReady value) {
if (!ready.isCompleted) ready.complete(value);
return;
}
worker?._handleResponse(message);
});
final StreamSubscription<Object?> errorSubscription = errors.listen((Object? message) {
final Object error = _remoteError(message);
if (!ready.isCompleted) {
ready.completeError(error);
} else if (worker case final NativeTaskIsolate activeWorker) {
activeWorker._handleTermination(error);
} else {
earlyTermination = error;
}
});
final StreamSubscription<Object?> exitSubscription = exits.listen((Object? _) {
final StateError error = StateError('Native task worker exited unexpectedly.');
if (!ready.isCompleted) {
ready.completeError(error);
} else if (worker case final NativeTaskIsolate activeWorker) {
activeWorker._handleTermination(error);
} else {
earlyTermination = error;
}
});
final Isolate isolate = await Isolate.spawn<_NativeTaskBootstrap>(
_runNativeTaskWorker,
_NativeTaskBootstrap(responses.sendPort, factory, initialMessage),
debugName: debugName,
onError: errors.sendPort,
onExit: exits.sendPort,
);
late final _NativeTaskReady handshake;
try {
handshake = await ready.future;
} on Object {
isolate.kill(priority: Isolate.immediate);
await responseSubscription.cancel();
await errorSubscription.cancel();
await exitSubscription.cancel();
responses.close();
errors.close();
exits.close();
rethrow;
}
if (handshake.error case final Object error) {
isolate.kill(priority: Isolate.immediate);
await responseSubscription.cancel();
await errorSubscription.cancel();
await exitSubscription.cancel();
responses.close();
errors.close();
exits.close();
Error.throwWithStackTrace(nativeTaskFailure(error), handshake.stackTrace ?? StackTrace.empty);
}
final NativeTaskIsolate result = NativeTaskIsolate._(
isolate,
responses,
responseSubscription,
<ReceivePort>[errors, exits],
<StreamSubscription<Object?>>[errorSubscription, exitSubscription],
handshake.commands!,
);
worker = result;
if (earlyTermination case final Object error) {
result._handleTermination(error);
}
return result;
}