stream<T> method
The typed events matching event. Broadcast, so several widgets may
listen to one channel.
Implementation
Stream<T> stream<T>(ChannelEvent<T> event) {
final channel = _channels.putIfAbsent(
event.channel,
() => _Channel(event.channel),
);
late final StreamController<T> controller;
StreamSubscription<RealtimeFrame>? frames;
controller = StreamController<T>.broadcast(
onListen: () {
channel.listeners++;
frames = channel.frames.stream.listen((frame) {
if (frame.event != event.event) return;
try {
controller.add(event.decode(frame.data));
} catch (error, stackTrace) {
// A payload this listener cannot read is this listener's problem.
// Every other channel on the socket keeps running.
controller.addError(error, stackTrace);
}
}, onDone: () => unawaited(controller.close()));
channel.errors.add(controller);
unawaited(_activate(channel));
},
onCancel: () async {
await frames?.cancel();
frames = null;
channel.errors.remove(controller);
if (--channel.listeners > 0) return;
await _deactivate(channel);
},
);
return controller.stream;
}