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. void ack(
    1. int dartPort
    )?,
  6. bool coalesced = false,
  7. String? debugLabel,
  8. @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 does NOT tear down the JS context — the old module instance and its emitters survive on the page; the bridge's nitro_web_instance_changed() ownership claim is what stands them down.

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,

  /// Coalesced streams (`Backpressure.batch` on an all-C++ spec): called
  /// after every delivered message so the bridge flushes what accumulated
  /// meanwhile, or goes idle.
  void Function(int dartPort)? ack,

  /// Each message is a `List` of items (the bridge batcher's delivery shape);
  /// [unpack] runs once per element, so a failing item forwards one error
  /// and the rest of the batch is still delivered.
  bool coalesced = false,
  String? debugLabel,
  @visibleForTesting WebReceivePort? testPort,
}) {
  // Built only when a line is actually logged: streams open per subscription.
  String label() => debugLabel ?? 'Stream<$T>';
  final receivePort = testPort ?? WebReceivePort();
  final nativePort = receivePort.sendPort.nativePort;
  var released = false;
  var eventCount = 0;

  if (_verboseOn) {
    _log(NitroLogLevel.verbose, label(), 'opening (port=$nativePort)');
  }

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

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

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

  void deliver(dynamic message) {
    try {
      final item = unpack(message);
      eventCount++;
      if (_verboseOn) {
        _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);
    }
  }

  receivePort.listen((dynamic message) {
    if (controller.isClosed) return;
    if (coalesced) {
      for (final m in message as List<dynamic>) {
        if (controller.isClosed) break;
        deliver(m);
      }
    } else {
      deliver(message);
    }
    if (!released) ack?.call(nativePort);
  });

  return controller.stream;
}