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.