Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
213 changes: 213 additions & 0 deletions src/DynamicData.Tests/List/SwitchFixture.cs
Original file line number Diff line number Diff line change
@@ -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;

Expand Down Expand Up @@ -60,4 +65,212 @@ public void PoulatesFirstSource()

inital.Should().BeEquivalentTo(_source.Items);
}

[Fact]
public void PropagatesOuterErrors()
{
using var source = new SourceList<int>();
using var switchable = new BehaviorSubject<IObservable<IChangeSet<int>>>(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<int>();
using var switchable = new BehaviorSubject<IObservable<IChangeSet<int>>>(source.Connect());
using var results = switchable.Switch().AsAggregator();

source.AddRange(Enumerable.Range(1, 100));

using var source2 = new Subject<IChangeSet<int>>();
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<int>();
using var switchable = new BehaviorSubject<IObservable<IChangeSet<int>>>(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<int>();
using var switchable = new BehaviorSubject<IObservable<IChangeSet<int>>>(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<int>();
using var second = new SourceList<int>();
using var switchable = new BehaviorSubject<IObservable<IChangeSet<int>>>(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<IChangeSet<int>>()).Switch().AsAggregator();

results.IsCompleted.Should().BeTrue("everything completed during subscription");
results.Exception.Should().BeNull();
}

[Fact]
public void DeliversChangesEmittedBeforeSynchronousCompletion()
{
var change = new ChangeSet<int> { new(ListChangeReason.Add, 42) };
using var results = Observable.Return(Observable.Return((IChangeSet<int>)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<IChangeSet<int>>();
using var second = new Subject<IChangeSet<int>>();
using var switchable = new BehaviorSubject<IObservable<IChangeSet<int>>>(first);
using var results = switchable.Switch().AsAggregator();

first.OnNext(new ChangeSet<int> { 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<int> { 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<IChangeSet<int>>(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<IObservable<IChangeSet<int>>>();
using var first = new Subject<IChangeSet<int>>();

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<int> { new Change<int>(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<IChangeSet<int>>())) { 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<IChangeSet<int>>);
var superseded = RawAnonymousObservable.Create<IChangeSet<int>>(observer =>
{
supersededObserver = observer;
return Disposable.Empty;
});

using var switchable = new Subject<IObservable<IChangeSet<int>>>();
using var current = new Subject<IChangeSet<int>>();

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<int> { new Change<int>(ListChangeReason.Add, 1, 0) });

results.Data.Count.Should().Be(1, "the selected source should still be delivering");
}
}
106 changes: 96 additions & 10 deletions src/DynamicData/List/Internal/Switch.cs
Original file line number Diff line number Diff line change
Expand Up @@ -15,21 +15,107 @@ internal sealed class Switch<T>(IObservable<IObservable<IChangeSet<T>>> sources)
public IObservable<IChangeSet<T>> Run() => Observable.Create<IChangeSet<T>>(
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<IChangeSet<T>>(observer);

var destination = new SourceList<T>();
// What the current source has contributed, so that switching away can take it back out.
var current = new List<T>();
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<T> { new Change<T>(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();
});
});
}
Loading