How do I prematurely terminate an observable chain using TakeUntil?

Viewed 133

I currently have the following:

var disconnect = Observable
    .FromEvent<ExceptionListener, Exception>(
        (handler) => connection.ExceptionListener += handler,
        (handler) => connection.ExceptionListener -= handler
    )
    .Do((exception) =>
    {
        // note: This line does get executed when I trigger the error scenario,
        //       so I know it is in fact happening.
        //
        Console.WriteLine(exception.Message);
    });

var messages = Observable
    .Using(
        () => connection.CreateSession(),
        (session) => Observable
    .Using(
        () => session.GetQueue(address),
        (queue) => Observable
    .Using(
        () => session.CreateConsumer(queue),
        (consumer) => Observable
            .FromEvent<MessageListener, IMessage>(
                (handler) => consumer.Listener += handler,
                (handler) => consumer.Listener -= handler
            )
            .TakeUntil(disconnect)
    )))

Whenever I build on or otherwise subscribe to messages, it seems like the TakeUntil is not being honoured and the chain isn't aborted, even though disconnect does in fact seem to emit.

Ideally I'd like for this IObservable (or any resultant one that builds on it) to complete/terminate (I'm not sure what the right word is here) as soon as the disconnect observable emits.

For completeness, my consuming code is effectively:

await messages.ForEachAsync(async (message) =>
{
    // note: Do things with `message`
}, cancellationToken);

// note: During normal operation, this never gets run. But I want the await 
//       above to complete when my `TakeUntil` above emits.
Console.WriteLine("Chain completed/terminated.");
2 Answers

Your code is working perfectly fine as it is. Here's how I tested it.

First up I got your original code into a compilable and runnable state.

public delegate void ExceptionListener();
public delegate void MessageListener();

public interface IMessage
{
    
}

public static class connection
{
    public static event ExceptionListener ExceptionListener;
    public static Session CreateSession() => new Session();
}

public class Session : IDisposable
{
    public Queue GetQueue(string address) => new Queue();
    public Consumer CreateConsumer(Queue queue) => new Consumer();
    
    #region IDisposable Support
    private bool disposedValue = false; // To detect redundant calls

    protected virtual void Dispose(bool disposing)
    {
        if (!disposedValue)
        {
            if (disposing)
            {
                // TODO: dispose managed state (managed objects).
            }

            // TODO: free unmanaged resources (unmanaged objects) and override a finalizer below.
            // TODO: set large fields to null.

            disposedValue = true;
        }
    }

    // TODO: override a finalizer only if Dispose(bool disposing) above has code to free unmanaged resources.
    // ~Session() {
    //   // Do not change this code. Put cleanup code in Dispose(bool disposing) above.
    //   Dispose(false);
    // }

    // This code added to correctly implement the disposable pattern.
    public void Dispose()
    {
        // Do not change this code. Put cleanup code in Dispose(bool disposing) above.
        Dispose(true);
        // TODO: uncomment the following line if the finalizer is overridden above.
        // GC.SuppressFinalize(this);
    }
    #endregion

}

public class Queue : IDisposable
{

    #region IDisposable Support
    private bool disposedValue = false; // To detect redundant calls

    protected virtual void Dispose(bool disposing)
    {
        if (!disposedValue)
        {
            if (disposing)
            {
                // TODO: dispose managed state (managed objects).
            }

            // TODO: free unmanaged resources (unmanaged objects) and override a finalizer below.
            // TODO: set large fields to null.

            disposedValue = true;
        }
    }

    // TODO: override a finalizer only if Dispose(bool disposing) above has code to free unmanaged resources.
    // ~Queue() {
    //   // Do not change this code. Put cleanup code in Dispose(bool disposing) above.
    //   Dispose(false);
    // }

    // This code added to correctly implement the disposable pattern.
    public void Dispose()
    {
        // Do not change this code. Put cleanup code in Dispose(bool disposing) above.
        Dispose(true);
        // TODO: uncomment the following line if the finalizer is overridden above.
        // GC.SuppressFinalize(this);
    }
    #endregion

}

public class Consumer : IDisposable
{
    public event MessageListener Listener;

    #region IDisposable Support
    private bool disposedValue = false; // To detect redundant calls

    protected virtual void Dispose(bool disposing)
    {
        if (!disposedValue)
        {
            if (disposing)
            {
                // TODO: dispose managed state (managed objects).
            }

            // TODO: free unmanaged resources (unmanaged objects) and override a finalizer below.
            // TODO: set large fields to null.

            disposedValue = true;
        }
    }

    // TODO: override a finalizer only if Dispose(bool disposing) above has code to free unmanaged resources.
    // ~Consumer() {
    //   // Do not change this code. Put cleanup code in Dispose(bool disposing) above.
    //   Dispose(false);
    // }

    // This code added to correctly implement the disposable pattern.
    public void Dispose()
    {
        // Do not change this code. Put cleanup code in Dispose(bool disposing) above.
        Dispose(true);
        // TODO: uncomment the following line if the finalizer is overridden above.
        // GC.SuppressFinalize(this);
    }
    #endregion
}

Now I can run this code:

var address = "";

var disconnect =
    Observable
        .FromEvent<ExceptionListener, Exception>(
            handler => connection.ExceptionListener += handler,
            handler => connection.ExceptionListener -= handler)
        .Do(exception => Console.WriteLine(exception.Message));

var messages =
    Observable.Using(
        () => connection.CreateSession(),
        session => Observable.Using(
            () => session.GetQueue(address),
            queue => Observable.Using(
                () => session.CreateConsumer(queue),
                consumer =>
                    Observable
                        .FromEvent<MessageListener, IMessage>(
                            handler => consumer.Listener += handler,
                            handler => consumer.Listener -= handler)
                        .TakeUntil(disconnect))));

Then I refactored to avoid the events and to use a couple of Subject<T> instances to mock the events:

var address = "";

var disconnectSubject = new Subject<Exception>();

var disconnect =
    disconnectSubject
        .Do(exception => Console.WriteLine(exception.Message));

var messageSubject = new Subject<IMessage>();

var messages =
    Observable.Using(
        () => connection.CreateSession(),
        session => Observable.Using(
            () => session.GetQueue(address),
            queue => Observable.Using(
                () => session.CreateConsumer(queue),
                consumer => messageSubject.TakeUntil(disconnect))));

Now I can run this code:

messages.Subscribe(m => Console.WriteLine("Message"));

messageSubject.OnNext(null);
messageSubject.OnNext(null);
messageSubject.OnNext(null);
disconnectSubject.OnNext(new Exception());
messageSubject.OnNext(null);

The output I get is:

Message
Message
Message
Exception of type 'System.Exception' was thrown.

The messages observable is indeed ending when the disconnect subject fires.

Then I replaced my simple messages.Subscribe code with this:

var cancellationToken = new CancellationToken();
var task =
    messages
        .ObserveOn(Scheduler.Default)
        .ForEachAsync(async (message) =>
        {
            await Task.Delay(TimeSpan.FromSeconds(1.0));
            Console.WriteLine("Message");
        }, cancellationToken);

I now get:

Exception of type 'System.Exception' was thrown.
Message
Message
Message

Apart from the change of output order produced by the await Task.Delay(TimeSpan.FromSeconds(1.0)); nothing has changed here.

The messages observable ends in both cases, as expected.

Your code, as presented here, works just fine.

The problem is probably with how the observable is consumed:

await messages.ForEachAsync(async (message) =>
{
    // note: Do things with `message`
});

The ForEachAsync operator does not accept an asynchronous delegate (Func<Task>), so the lambda passed is async void (also known as fire-and-crash). If you want to invoke an asynchronous lambda for each element of an observable sequence, you must use a projection to create an IObservable<Task<T>>, and then use one of the three flattening operators (Concat, Merge or Switch) to get back an IObservable<TResult> containing the results of the asynchronous invocations. If these don't have results, you can return something dummy like the Unit.Default.

await messages.Select(async message =>
{
    // note: Do things with `message`
    return Unit.Default;
}).Merge().DefaultIfEmpty();

The DefaultIfEmpty is there just to prevent an InvalidOperationException in case the sequence completes with zero messages.


Example with cancelable lambda:

await messages.Select(message => Observable.FromAsync(async cancellationToken =>
{
    // note: Do things with `message`, while observing the cancellationToken
    return Unit.Default;
})).Merge().DefaultIfEmpty();
Related