stream<T> method

Stream<T> stream<T>(
  1. ChannelEvent<T> event
)

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;
}