openStream<T> static method

Stream<T> openStream<T>({
  1. required void register(
    1. int dartPort
    ),
  2. required T unpack(
    1. dynamic message
    ),
  3. required void release(
    1. int dartPort
    ),
  4. required Backpressure backpressure,
  5. String? debugLabel,
  6. @visibleForTesting WebReceivePort? testPort,
})

Opens a stream from a WASM event source over a WebReceivePort, mirroring the native lifecycle (explicit cancel, GC finalizer safety net; hot restart tears down the JS context wholesale).

Implementation

static Stream<T> openStream<T>({
  required void Function(int dartPort) register,
  required T Function(dynamic message) unpack,
  required void Function(int dartPort) release,
  required Backpressure backpressure,
  String? debugLabel,
  @visibleForTesting WebReceivePort? testPort,
}) {
  final label = debugLabel ?? 'Stream<$T>';
  final receivePort = testPort ?? WebReceivePort();
  final nativePort = receivePort.sendPort.nativePort;
  var released = false;
  var eventCount = 0;

  _log(NitroLogLevel.verbose, label, 'opening (port=$nativePort)');

  void doRelease() {
    if (released) return;
    released = true;
    _log(
      NitroLogLevel.verbose,
      label,
      'releasing (port=$nativePort, events=$eventCount)',
    );
    release(nativePort);
    receivePort.close();
  }

  final controller = StreamController<T>(
    onListen: () {
      _log(NitroLogLevel.verbose, label, 'listener attached — registering');
      register(nativePort);
    },
    onCancel: doRelease,
  );

  _streamFinalizer.attach(controller, doRelease, detach: controller);

  receivePort.listen((dynamic message) {
    if (controller.isClosed) return;
    try {
      final item = unpack(message);
      eventCount++;
      _log(
        NitroLogLevel.verbose,
        label,
        'event #$eventCount unpacked',
      );
      controller.add(item);
    } catch (e, st) {
      _log(
        NitroLogLevel.error,
        label,
        'unpack failed on event #${eventCount + 1} — forwarding error to stream',
        e,
        st,
      );
      controller.addError(e, st);
    }
  });

  return controller.stream;
}