openStream<T> static method
Stream<T>
openStream<T>({
- required void register(
- int dartPort
- required T unpack(
- dynamic message
- required void release(
- int dartPort
- required Backpressure backpressure,
- void ack(
- int dartPort
- bool coalesced = false,
- String? debugLabel,
- @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;
}