I have the following setup

IObservable<Data> source = ...;

source
    .Select(data=>VeryExpensiveOperation(data))
    .Subscribe(data=>Console.WriteLine(data));

Normally the events come seperated by a reasonable time frame. Imagine a user updating a text box in a form. Our VeryExpensiveOperation might take 5 seconds to complete and whilst it does an hour glass is displayed on the screen.

However if during the 5 seconds the user updates the textbox again I would want to send a cancelation to the current VeryExpensiveOperation before the new one starts.

I would imagine a scenario like

source
    .SelectWithCancel((data, cancelToken)=>VeryExpensiveOperation(data, token))
    .Subscribe(data=>Console.WriteLine(data));

So every time the lambda is called is is called with a cancelToken which can be used to manage canceling a Task. However now we are mixing Task, CancelationToken and RX. Not quite sure how to fit it all together. Any suggestions.

Bonus Points for figuring out how to test the operator using XUnit :)

FIRST ATTEMPT

    public static IObservable<U> SelectWithCancelation<T, U>( this IObservable<T> This, Func<CancellationToken, T, Task<U>> fn )
    {
        CancellationTokenSource tokenSource = new CancellationTokenSource();

        return This
            .ObserveOn(Scheduler.Default)
            .Select(v=>{
                tokenSource.Cancel();
                tokenSource=new CancellationTokenSource();
                return new {tokenSource.Token, v};
            })
            .SelectMany(o=>Observable.FromAsync(()=>fn(o.Token, o.v)));
    }

Not tested yet. I'm hoping that a task that does not complete generates an IObservable that completes without firing any OnNext events.

Edit
Report