diff --git a/src/DynamicData.Tests/List/SwitchFixture.cs b/src/DynamicData.Tests/List/SwitchFixture.cs index abcdaa8d3..5b490f84f 100644 --- a/src/DynamicData.Tests/List/SwitchFixture.cs +++ b/src/DynamicData.Tests/List/SwitchFixture.cs @@ -1,6 +1,11 @@ using System; using System.Linq; +using System.Reactive.Disposables; +using System.Reactive.Linq; using System.Reactive.Subjects; +using System.Threading; + +using DynamicData.Tests.Utilities; using FluentAssertions; @@ -60,4 +65,212 @@ public void PoulatesFirstSource() inital.Should().BeEquivalentTo(_source.Items); } + + [Fact] + public void PropagatesOuterErrors() + { + using var source = new SourceList(); + using var switchable = new BehaviorSubject>>(source.Connect()); + using var results = switchable.Switch().AsAggregator(); + + source.AddRange(Enumerable.Range(1, 100)); + + var error = new Exception("Test"); + switchable.OnError(error); + + results.Exception.Should().Be(error); + } + + [Fact] + public void PropagatesInnerErrors() + { + using var source = new SourceList(); + using var switchable = new BehaviorSubject>>(source.Connect()); + using var results = switchable.Switch().AsAggregator(); + + source.AddRange(Enumerable.Range(1, 100)); + + using var source2 = new Subject>(); + switchable.OnNext(source2); + + var error = new Exception("Test"); + source2.OnError(error); + + results.Exception.Should().Be(error); + } + + [Fact] + public void CompletesWhenSourcesAndInnerComplete() + { + using var source = new SourceList(); + using var switchable = new BehaviorSubject>>(source.Connect()); + using var results = switchable.Switch().AsAggregator(); + + source.AddRange(Enumerable.Range(1, 100)); + + switchable.OnCompleted(); + results.IsCompleted.Should().BeFalse("the inner sequence is still running"); + + source.Dispose(); + + results.IsCompleted.Should().BeTrue("both the sources and the inner sequence have completed"); + results.Exception.Should().BeNull(); + results.Data.Count.Should().Be(100, "all data should have been received before completion"); + } + + [Fact] + public void DoesNotCompleteWhileInnerIsStillRunning() + { + using var source = new SourceList(); + using var switchable = new BehaviorSubject>>(source.Connect()); + using var results = switchable.Switch().AsAggregator(); + + switchable.OnCompleted(); + source.Add(1); + + results.IsCompleted.Should().BeFalse("the inner sequence has not completed"); + results.Data.Count.Should().Be(1, "changes should still flow after the sources sequence completes"); + } + + [Fact] + public void DoesNotCompleteWhenOnlyASupersededInnerCompletes() + { + using var first = new SourceList(); + using var second = new SourceList(); + using var switchable = new BehaviorSubject>>(first.Connect()); + using var results = switchable.Switch().AsAggregator(); + + switchable.OnNext(second.Connect()); + switchable.OnCompleted(); + + first.Dispose(); + + results.IsCompleted.Should().BeFalse("the superseded sequence is not the current one"); + + second.Add(1); + results.Data.Count.Should().Be(1, "the current sequence should still be delivering"); + + second.Dispose(); + results.IsCompleted.Should().BeTrue("the current sequence has now completed"); + } + + [Fact] + public void CompletesWhenSourcesAndInnerCompleteSynchronously() + { + using var results = Observable.Return(Observable.Empty>()).Switch().AsAggregator(); + + results.IsCompleted.Should().BeTrue("everything completed during subscription"); + results.Exception.Should().BeNull(); + } + + [Fact] + public void DeliversChangesEmittedBeforeSynchronousCompletion() + { + var change = new ChangeSet { new(ListChangeReason.Add, 42) }; + using var results = Observable.Return(Observable.Return((IChangeSet)change)).Switch().AsAggregator(); + + results.Data.Count.Should().Be(1, "changes emitted before a synchronous completion must not be lost"); + results.IsCompleted.Should().BeTrue("the source completed"); + results.Exception.Should().BeNull(); + } + + [Fact] + public void IgnoresChangesFromASupersededSource() + { + using var first = new Subject>(); + using var second = new Subject>(); + using var switchable = new BehaviorSubject>>(first); + using var results = switchable.Switch().AsAggregator(); + + first.OnNext(new ChangeSet { new(ListChangeReason.Add, 1) }); + results.Data.Count.Should().Be(1); + + switchable.OnNext(second); + results.Data.Count.Should().Be(0, "moving to a new source drops what the previous one contributed"); + + first.OnNext(new ChangeSet { new(ListChangeReason.Add, 2) }); + + results.Data.Count.Should().Be(0, "a superseded source must not be able to write into the result"); + results.Exception.Should().BeNull(); + } + + [Fact] + public void PropagatesInnerErrorsRaisedSynchronously() + { + var error = new Exception("Test"); + using var results = Observable.Return(Observable.Throw>(error)).Switch().AsAggregator(); + + results.Exception.Should().Be(error, "the error was raised during subscription"); + } + + [Fact] + public void DoesNotHoldALockWhileDeliveringDownstream() + { + // Observable.Switch holds its gate for the whole of downstream delivery, which is the shape that + // deadlocks when a pipeline crosses into another collection. Delivery has to go through the queue, + // which enqueues and returns, so a producer is never held up by whatever a subscriber is doing. + using var switchable = new Subject>>(); + using var first = new Subject>(); + + using var isDelivering = new ManualResetEventSlim(false); + using var release = new ManualResetEventSlim(false); + + using var subscription = switchable.Switch().Subscribe(_ => + { + isDelivering.Set(); + release.Wait(TimeSpan.FromSeconds(10)); + }); + + switchable.OnNext(first); + + var deliverer = new Thread(() => first.OnNext(new ChangeSet { new Change(ListChangeReason.Add, 1, 0) })) { IsBackground = true }; + deliverer.Start(); + + isDelivering.Wait(TimeSpan.FromSeconds(10)).Should().BeTrue("the subscriber should have been handed the change"); + + var producer = new Thread(() => switchable.OnNext(new Subject>())) { IsBackground = true }; + producer.Start(); + + var producerFinished = producer.Join(TimeSpan.FromSeconds(2)); + + release.Set(); + deliverer.Join(TimeSpan.FromSeconds(10)); + producer.Join(TimeSpan.FromSeconds(10)); + + producerFinished.Should().BeTrue("writing to the source must not block while a subscriber holds onto a notification"); + } + + [Fact] + public void IgnoresErrorsFromASupersededSource() + { + // Switching away from a source means anything it produces afterwards belongs to a source that + // is no longer selected, and that includes its failures. Ordinarily disposal stops a + // superseded source being heard from again, but disposal cannot reach a notification that is + // already in flight, so the operator has to discard it on arrival. The raw observable hands + // back the observer directly, which is how that in-flight failure is reproduced here without + // needing a race to land. + var supersededObserver = default(IObserver>); + var superseded = RawAnonymousObservable.Create>(observer => + { + supersededObserver = observer; + return Disposable.Empty; + }); + + using var switchable = new Subject>>(); + using var current = new Subject>(); + + using var results = switchable.Switch().AsAggregator(); + + switchable.OnNext(superseded); + switchable.OnNext(current); + + supersededObserver.Should().NotBeNull("the superseded source should have been subscribed"); + supersededObserver!.OnError(new Exception("Test")); + + results.Exception.Should().BeNull("the failed source had already been switched away from"); + + current.OnNext(new ChangeSet { new Change(ListChangeReason.Add, 1, 0) }); + + results.Data.Count.Should().Be(1, "the selected source should still be delivering"); + } } diff --git a/src/DynamicData/List/Internal/Switch.cs b/src/DynamicData/List/Internal/Switch.cs index 9425771dd..dcd20c283 100644 --- a/src/DynamicData/List/Internal/Switch.cs +++ b/src/DynamicData/List/Internal/Switch.cs @@ -15,21 +15,107 @@ internal sealed class Switch(IObservable>> sources) public IObservable> Run() => Observable.Create>( observer => { - var locker = InternalEx.NewLock(); + // Switching is done by hand rather than with Observable.Switch, which holds its gate for + // the whole of downstream delivery. The queue enqueues and returns instead, so a producer + // is never held up by whatever a subscriber does with the notification, and a pipeline + // crossing into another collection cannot deadlock against it. + var queue = new DeliveryQueue>(observer); - var destination = new SourceList(); + // What the current source has contributed, so that switching away can take it back out. + var current = new List(); + var subscription = new SerialDisposable(); - var populator = Observable.Switch( - _sources.Do( - _ => + // Identifies the current source. A superseded one may still be mid-delivery, and anything + // it produces after this point belongs to a source that has already been switched away from. + var active = 0; + var isSourceRunning = false; + var areSourcesComplete = false; + + var outer = _sources.SubscribeSafe( + source => + { + int id; + + using (var scope = queue.AcquireLock()) { - lock (locker) + id = ++active; + isSourceRunning = true; + + if (current.Count != 0) { - destination.Clear(); + scope.EnqueueNext(new ChangeSet { new Change(ListChangeReason.Clear, current.ToArray()) }); + current.Clear(); } - })).Synchronize(locker).PopulateInto(destination); + } + + // Subscribed outside the lock. The source may deliver synchronously, and that + // delivery takes the lock for itself. + subscription.Disposable = source.SubscribeSafe( + changes => + { + using var scope = queue.AcquireLock(); + + if (id != active) + { + return; + } + + current.Clone(changes); + + if (changes.Count != 0) + { + scope.EnqueueNext(changes); + } + }, + error => + { + using var scope = queue.AcquireLock(); + + if (id != active) + { + return; + } + + scope.EnqueueError(error); + }, + () => + { + using var scope = queue.AcquireLock(); + + if (id != active) + { + return; + } + + isSourceRunning = false; + + if (areSourcesComplete) + { + scope.EnqueueCompleted(); + } + }); + }, + queue.OnError, + () => + { + using var scope = queue.AcquireLock(); + + areSourcesComplete = true; + + // The current source may still be running, and the result ends only once both have. + if (!isSourceRunning) + { + scope.EnqueueCompleted(); + } + }); - var publisher = destination.Connect().SubscribeSafe(observer); - return new CompositeDisposable(destination, populator, publisher); + // Disposal order matters and CompositeDisposable does not specify one. The queue goes first + // so that any delivery in flight is finished before the subscriptions feeding it are torn down. + return Disposable.Create(() => + { + queue.Dispose(); + outer.Dispose(); + subscription.Dispose(); + }); }); }