subscribe method

AcpOutboundSubscription<T> subscribe()

Creates an independent live subscription.

Implementation

AcpOutboundSubscription<T> subscribe() {
  if (_isClosed) {
    return AcpOutboundSubscription<T>._(
      replay: List<T>.unmodifiable(<T>[]),
      live: Stream<T>.empty(),
    );
  }

  final List<T> replay;
  if (_hasSubscribed) {
    replay = List<T>.unmodifiable(<T>[]);
  } else {
    _hasSubscribed = true;
    replay = List<T>.unmodifiable(_replay);
    _replay.clear();
  }

  late final _HubSubscriber<T> subscriber;
  subscriber = _HubSubscriber<T>(
    capacity: capacity,
    onOverflow: onOverflow,
    onCancel: () {
      _subscribers.remove(subscriber);
    },
  );
  _subscribers.add(subscriber);
  return AcpOutboundSubscription<T>._(
    replay: replay,
    live: subscriber.stream,
  );
}