spawn static method

Future<NativeTaskIsolate> spawn({
  1. required NativeTaskWorkerFactory factory,
  2. required Object? initialMessage,
  3. 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;
}