So I couldn't find any way to do it with operators, but we can always create our own connectable observable class! :)
public class UntilSubscribeReplay<T> : IObservable<T>
{
private readonly IObservable<T> _parent;
private HashSet<IObserver<T>> _observers = new HashSe
private IDisposable _subscription;
private readonly Queue<T> _itemsQueue = new Queue<T>(
public UntilSubscribeReplay(IObservable<T> parent)
{
_parent = parent;
Subscribe();
}
public IDisposable Subscribe(IObserver<T> observer)
{
if (_subscription == null)
Subscribe();
_observers.Add(observer);
EmitQueuedItems();
return Disposable.Create(() =>
{
_observers.Remove(observer);
if (_observers.Count == 0)
_subscription.Dispose();
});
}
private void Subscribe()
{
_subscription = _parent
.Do(_itemsQueue.Enqueue)
.Subscribe(_ => EmitQueuedItems());
}
private void EmitQueuedItems()
{
if (_observers.Count == 0)
return;
while (_itemsQueue.TryDequeue(out T item))
{
foreach (IObserver<T> observer in _observers)
{
observer.OnNext(item);
}
}
}
}
And of course, an extension method to use for your convenience. :)
public static class MyObservableExtensions
{
public static IObservable<T> UntilSubscribeReplay<T>(this IObservable<T> source)
=> new UntilSubscribeReplay<T>(source);
}
And here is how you would use it:
IObservable<int> myReplay = myObservable.UntilSubscribeReplay();
// Do Stuff
myReplay.Subscribe(value => HandleValue(value));
Note, that there is no need to use the connect. The proper way would be to implement it as a Connectable Observable (To avoid memory leaks when you use the extension but don't subscribe eventually), but for your specific case, it should be enough.
Let me know how it goes. :)