Strange Rx+CancellationToken issue: sometimes the registered callback does not complete

Viewed 109

I observed a strange phenomenon that occurs sometimes with an Rx query I wrote, that involves a CancellationToken. Two callbacks are registered to the same CancellationToken, one outside of the query and one that is part of the query. The intention of the CancellationToken is to signal the termination of the query. What happens is that sometimes the second callback is stuck in the middle of the execution, never completing, preventing the first callback from being invoked.

Below is a minimal example that reproduces the issue. It's not very minimal, but I can't reduce it any further. For example replacing the Switch operator with the Merge makes the issue disappear. The same happens if the exception thrown by the Task.Delay(1000, cts.Token) is swallowed.

public class Program
{
    public static void Main()
    {
        var cts = new CancellationTokenSource(500);
        cts.Token.Register(() => Console.WriteLine("### Token Canceled! ###"));
        try
        {
            Observable
                .Timer(TimeSpan.Zero, TimeSpan.FromMilliseconds(1000))
                .TakeUntil(Observable.Create<Unit>(observer =>
                    cts.Token.Register(() =>
                    {
                        Console.WriteLine("Before observer.OnNext");
                        observer.OnNext(Unit.Default);
                        Console.WriteLine("After observer.OnNext");
                    })))
                .Select(_ =>
                {
                    return Observable.StartAsync(async () =>
                    {
                        Console.WriteLine("Action starting");
                        await Task.Delay(1000, cts.Token);
                        return 1;
                    });
                })
                .Switch()
                .Wait();
        }
        catch (Exception ex) { Console.WriteLine("Failed: {0}", ex.Message); }
        Thread.Sleep(500);
        Console.WriteLine("Finished");
    }
}

Expected output:

Action starting
Before observer.OnNext
After observer.OnNext
### Token Canceled! ###
Failed: A task was canceled.
Finished

Actual output (sometimes):

Action starting
Before observer.OnNext
Failed: A task was canceled.
Finished

Try it on fiddle. You may need to run the program 3-4 times before the issue appears. Notice the two missing log entries. It seems that the call observer.OnNext(Unit.Default); never completes.

My question is: Does anyone have any idea what causes this issue? Also, how could I modify the CancellationToken-related part of the query, so that it performs its intended purpose (terminates the query), without interfering with other registered callbacks of the same CancellationToken?

.NET 5.0.1 & .NET Framework 4.8, System.Reactive 5.0.0, C# 9

Update: also .NET 6.0 with System.Reactive 5.0.0 (screenshot taken at June 4, 2022)


One more observation: The issue stops appearing if I modify the Observable.Create delegate so that it returns a Disposable.Empty instead of a CancellationTokenRegistration, like this:

.TakeUntil(Observable.Create<Unit>(observer =>
{
    cts.Token.Register(() =>
    {
        Console.WriteLine("Before observer.OnNext");
        observer.OnNext(default);
        Console.WriteLine("After observer.OnNext");
    });
    return Disposable.Empty;
}))

But I don't think that ignoring the registration returned by the cts.Token.Register is a fix.

2 Answers

1yr later, I can't repro it. Are you still able to? I've tried it with .NET Framework 4.8 and .NET 6.

Yet, in my related question, I've just dealt with a race condition resulting in a deadlock, which might be related to the issue you've described.

CancellationTokenRegistration.Dispose is a blocking call by its nature, it waits for all callbacks to complete, and those callbacks are executed on the thread that called CancellationTokenSource.Cancel (e.g., a random timer callback thread), which is likely different from a thread that invokes Dispose of the Rx subscription.

I haven't dug deep, but the workaround in my case was to not invoke the token registration's Dispose if the token is already cancelled. This should be OK, as the token registrations are discarded after the cancellation is signaled, and we're not interested in the cancellation callback anymore when Rx calls Dispose:

    .TakeUntil(Observable.Create<Unit>(observer => 
    {
        var token = cts.Token;
        var rego = token.Register(() =>
        {
            Console.WriteLine("Before observer.OnNext");
            observer.OnNext(Unit.Default);
            Console.WriteLine("After observer.OnNext");
        });
        return Disposable.Create(() => 
        {
            if (!token.IsCancellationRequested)
            {
                // this 
                rego.Dispose();
            }
        });
    }))       

In .NET 6, one other possible workaround might be to use CancellationTokenRegistration.DisposeAsync (though I haven't verified if it would cure the deadlock in my case):

        return Disposable.Create(() => 
        {
            DisposeAsync();
            async void DisposeAsync() => 
              await rego.DisposeAsync().ConfigureAwait(false);
        });

Updated, as Theodor mentioned in the comments, using CancellationTokenRegistration.Unregister instead of Dispose is yet another workaround (non-blocking), which feels even cleaner to me.

After @noseratio posted his workaround, I found another one. I replaced the call CancellationTokenRegistration.Dispose with the call CancellationTokenRegistration.Unregister, and I can no longer reproduce the issue:

var cts = new CancellationTokenSource(500);
cts.Token.Register(() => Console.WriteLine("### Token Canceled! ###"));
try
{
    Observable
        .Timer(TimeSpan.Zero, TimeSpan.FromMilliseconds(1000))
        .TakeUntil(Observable.Create<Unit>(observer =>
            Disposable.Create(cts.Token.Register(() =>
            {
                Console.WriteLine("Before observer.OnNext");
                observer.OnNext(Unit.Default);
                Console.WriteLine("After observer.OnNext");
            }), state => state.Unregister())))
        .Select(_ =>
        {
            return Observable.StartAsync(async () =>
            {
                Console.WriteLine("Action starting");
                await Task.Delay(1000, cts.Token);
                return 1;
            });
        })
        .Switch()
        .Wait();
}
catch (Exception ex) { Console.WriteLine("Failed: {0}", ex.Message); }
Thread.Sleep(500);
Console.WriteLine("Finished");

Output (always 6 entries, as expected):

Action starting
Before observer.OnNext
After observer.OnNext
### Token Canceled! ###
Failed: A task was canceled.
Finished

Try it on Fiddle.

Related