subscribe method
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,
);
}