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
34 changes: 27 additions & 7 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -1142,6 +1142,7 @@ share one subscription instead of starting the work over each time.
| `IObserver<T>.FastForEach(source)` | Pushes a whole collection into an observer and indexes arrays and lists directly. | - |
| `ObserveOn(sequencer)` | Delivers notifications to subscribers on the sequencer you name. | `ObserveOn` |
| `WitnessOn(sequencer)` | Another name for `ObserveOn`. | `ObserveOn` |
| `WitnessLatestOn(sequencer)` | Delivers notifications on the sequencer you name. Keeps only the newest value that is waiting and drops the older ones. | - |
| `SubscribeOn(sequencer)` | Runs the subscription itself on the sequencer you name. | `SubscribeOn` |
| `Synchronize()`, `Synchronize(object gate)` | Delivers notifications one at a time behind a lock, which you can share with other signals. | `Synchronize` |
| `Synchronize(Lock gate)` | The same, taking a `System.Threading.Lock`. Available on net9.0 and later only. | `Synchronize` |
Expand Down Expand Up @@ -2211,18 +2212,36 @@ conveniences that take an action, with or without a delay.

| Sequencer | What it does |
|---|---|
| `Sequencer.CurrentThread` | Queues work on the calling thread and runs it in order. Also available as `CurrentThreadSequencer.Instance`. |
| `Sequencer.Immediate` | Runs the work right away on the calling thread. Also available as `ImmediateSequencer.Instance`. |
| `Sequencer.CurrentThread` | Queues work on the calling thread and runs it in order. Also available as `CurrentThreadSequencer.Instance`. `new CurrentThreadSequencer(timeProvider)` waits through a `TimeProvider`. |
| `Sequencer.Immediate` | Runs the work right away on the calling thread. Also available as `ImmediateSequencer.Instance`. `new ImmediateSequencer(timeProvider)` waits through a `TimeProvider`. |
| `Sequencer.Default` | The default choice for background work. It is `TaskPoolSequencer.Default`. |
| `TaskPoolSequencer` | Runs work through a `TaskFactory`. Use `TaskPoolSequencer.Instance`, or pass your own factory. |
| `ThreadPoolSequencer` | Runs work on the thread pool. Use `ThreadPoolSequencer.Instance`. |
| `SynchronizationContextSequencer` | Posts work to a `SynchronizationContext`. Use `SynchronizationContextSequencer.Current`, or pass a context. |
| `WasmSequencer` | Runs work on a browser's single-threaded event loop. Use `WasmSequencer.Default`. |
| `VirtualClock` | Controls time in a test, using `DateTimeOffset` and `TimeSpan`. You advance the clock, so nothing sleeps. |
| `TaskPoolSequencer` | Runs work through a `TaskFactory`. Use `TaskPoolSequencer.Instance`, or pass your own factory. `new TaskPoolSequencer(factory, timeProvider)` takes its time from a `TimeProvider`. |
| `ThreadPoolSequencer` | Runs work on the thread pool. Use `ThreadPoolSequencer.Instance`. `new ThreadPoolSequencer(timeProvider)` takes its time and delay timer from a `TimeProvider`. |
| `SynchronizationContextSequencer` | Posts work to a `SynchronizationContext`. Use `SynchronizationContextSequencer.Current`, or pass a context. `new SynchronizationContextSequencer(context, timeProvider)` takes its time from a `TimeProvider`. |
| `WasmSequencer` | Runs work on a browser's single-threaded event loop. Use `WasmSequencer.Default`. `new WasmSequencer(timeProvider)` takes its time and timers from a `TimeProvider`. |
| `VirtualClock` | Controls time in a test, using `DateTimeOffset` and `TimeSpan`. You advance the clock, so nothing sleeps. It is also a `TimeProvider`. |
| `VirtualTimeSequencer<TAbsolute, TRelative>` | Controls time in a test using your own clock types. |
| `sequencer.IsImmediate` | Returns `true` when the sequencer is the immediate sequencer and `false` for any other sequencer or for `null`. Use it instead of comparing against `Sequencer.Immediate`. It comes from `SequencerImmediacyExtensions`. |

Use a virtual sequencer for a time-based test. Do not sleep a real thread.

To run a built-in sequencer from a fake clock, pass a `TimeProvider` to its constructor. `VirtualClock` is a
`TimeProvider`, so `clock.AdvanceBy` fires the sequencer's timers.

```csharp
var clock = new VirtualClock();
var sequencer = new ThreadPoolSequencer(clock);
var results = new List<int>();

using var subscription = source
.Throttle(TimeSpan.FromSeconds(1), sequencer)
.Subscribe(results.Add);

source.OnNext(1);
clock.AdvanceBy(TimeSpan.FromSeconds(1));
// results now holds 1
```

### Platform sequencers

Each UI framework gets its own package, so the core package pulls in none of them.
Expand Down Expand Up @@ -3042,6 +3061,7 @@ the same thing.
| `DeferSignal<T>` | `Signal.Lazy` / `Signal.Defer` |
| `CatchSignal<T>` | `Catch` |
| `WitnessOnSignal<T>` | `WitnessOn` / `ObserveOn` |
| `WitnessLatestOnSignal<T>` | `WitnessLatestOn` |
| `SelectManyThenCoordinator<TSource, TMid, TResult>` | `SelectManyThen`, both projection stages in one sink |
| `CalmCoordinator<T>` | `Calm` / `Throttle` |
| `ReattemptCoordinator<T>` | `Reattempt` |
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ public IDisposable Subscribe(IObserver<T> observer)
InvalidOperationExceptionHelper.ThrowIfNull(scheduler);
ArgumentExceptionHelper.ThrowIfNull(observer);

if (ReferenceEquals(scheduler, Sequencer.Immediate))
if (scheduler.IsImmediate)
{
return source.Subscribe(observer);
}
Expand Down
2 changes: 1 addition & 1 deletion src/Primitives.Shared/Advanced/EmptySignal{T}.cs
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,7 @@ private IDisposable SubscribeCore(IObserver<T> observer, IDisposable cancel)
{
observer = new GuardedWitness<T>(observer, cancel);

if (_scheduler == Sequencer.Immediate)
if (_scheduler.IsImmediate)
{
observer.OnCompleted();
return EmptyDisposable.Instance;
Expand Down
2 changes: 1 addition & 1 deletion src/Primitives.Shared/Advanced/ReturnSignal{T}.cs
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,7 @@ private IDisposable SubscribeCore(IObserver<T> observer, IDisposable cancel)
{
observer = new GuardedWitness<T>(observer, cancel);

if (_scheduler == Sequencer.Immediate)
if (_scheduler.IsImmediate)
{
observer.OnNext(_value);
observer.OnCompleted();
Expand Down
2 changes: 1 addition & 1 deletion src/Primitives.Shared/Advanced/SubscriptionScheduling.cs
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,7 @@ internal static IDisposable OnCurrentThread<TState>(TState state, Func<TState, I
/// <returns>A disposable that cancels the scheduled work, best effort.</returns>
internal static IDisposable RunOn<TState>(ISequencer sequencer, TState state, Action<TState> run)
{
if (sequencer == Sequencer.Immediate)
if (sequencer.IsImmediate)
{
run(state);
return EmptyDisposable.Instance;
Expand Down
2 changes: 1 addition & 1 deletion src/Primitives.Shared/Advanced/ThrowSignal{T}.cs
Original file line number Diff line number Diff line change
Expand Up @@ -64,7 +64,7 @@ private IDisposable SubscribeCore(IObserver<T> observer, IDisposable cancel)
{
observer = new GuardedWitness<T>(observer, cancel);

if (_scheduler == Sequencer.Immediate)
if (_scheduler.IsImmediate)
{
observer.OnError(_error);
return EmptyDisposable.Instance;
Expand Down
244 changes: 244 additions & 0 deletions src/Primitives.Shared/Advanced/WitnessLatestOnSignal{T}.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,244 @@
// Copyright (c) 2019-2026 ReactiveUI Association Incorporated. All rights reserved.
// ReactiveUI Association Incorporated licenses this file to you under the MIT license.
// See the LICENSE file in the project root for full license information.

using System.Runtime.CompilerServices;

#if REACTIVE_SHIM
namespace ReactiveUI.Primitives.Reactive.Advanced;
#else
namespace ReactiveUI.Primitives.Advanced;
#endif

/// <summary>Re-dispatches source notifications onto a sequencer, delivering only the newest value that is waiting.</summary>
/// <typeparam name="T">The value type.</typeparam>
/// <param name="source">The source observable.</param>
/// <param name="scheduler">The sequencer that dispatches the notifications.</param>
[System.Diagnostics.DebuggerDisplay("WitnessLatestOnSignal<{typeof(T).Name,nq}>")]
public sealed class WitnessLatestOnSignal<T>(IObservable<T> source, ISequencer scheduler) : IRequireCurrentThread<T>
{
/// <summary>The source observable.</summary>
private readonly IObservable<T> _source = source;

/// <summary>The sequencer that dispatches the notifications.</summary>
private readonly ISequencer _scheduler = scheduler;

/// <summary>Reports that subscription is dispatched through the current-thread sequencer.</summary>
/// <returns>Always <see langword="true"/>.</returns>
[MethodImpl(MethodImplOptions.AggressiveInlining)]
public bool IsRequiredSubscribeOnCurrentThread() => true;

/// <summary>Subscribes an observer through the current-thread sequencer.</summary>
/// <param name="observer">The downstream observer.</param>
/// <returns>The subscription.</returns>
[MethodImpl(MethodImplOptions.AggressiveInlining)]
public IDisposable Subscribe(IObserver<T> observer) =>
SignalSubscription.Subscribe(observer, true, SubscribeCore);

/// <summary>Builds the dispatching sink and subscribes it to the source.</summary>
/// <param name="observer">The downstream observer.</param>
/// <param name="cancel">The subscription handle owned by the subscription helper.</param>
/// <returns>The source subscription together with the sink.</returns>
[MethodImpl(MethodImplOptions.AggressiveInlining)]
private IDisposable SubscribeCore(IObserver<T> observer, IDisposable cancel) =>
new WitnessLatestOn(this, observer, cancel).Run();

/// <summary>Holds the newest source value and delivers it on the sequencer.</summary>
/// <param name="parent">The owning signal.</param>
/// <param name="observer">The downstream observer.</param>
/// <param name="cancel">The subscription handle released on teardown.</param>
private sealed class WitnessLatestOn(WitnessLatestOnSignal<T> parent, IObserver<T> observer, IDisposable cancel) : IObserver<T>, IWorkItem, IsDisposed
{
/// <summary>The owning signal.</summary>
private readonly WitnessLatestOnSignal<T> _parent = parent;

/// <summary>The downstream observer.</summary>
private readonly IObserver<T> _observer = observer;

/// <summary>Synchronization gate guarding the pending slot and scheduling state.</summary>
private readonly Lock _gate = new();

/// <summary>Upstream subscription disposed on teardown.</summary>
private readonly IDisposable _cancel = cancel;

/// <summary>The newest value awaiting delivery.</summary>
private T _pending = default!;

/// <summary>The error awaiting delivery.</summary>
private Exception? _error;

/// <summary>Whether <see cref="_pending"/> holds a value.</summary>
private bool _hasPending;

/// <summary>Whether completion is awaiting delivery.</summary>
private bool _hasCompleted;

/// <summary>Whether the sink has been torn down.</summary>
private bool _isDisposed;

/// <summary>Tracks whether a drain is scheduled or running.</summary>
private bool _isScheduled;

/// <inheritdoc/>
public bool IsDisposed
{
get
{
lock (_gate)
{
return _isDisposed;
}
}
}

/// <summary>Subscribes to the source.</summary>
/// <returns>The source subscription together with this sink.</returns>
public MultipleDisposable Run()
{
var sourceDisposable = _parent._source.Subscribe(this);

return new(sourceDisposable, this);
}

/// <summary>Stores a value as the newest pending value.</summary>
/// <param name="value">The value to store.</param>
public void OnNext(T value)
{
lock (_gate)
{
if (_isDisposed || _hasCompleted || _error is not null)
{
return;
}

_pending = value;
_hasPending = true;
if (_isScheduled)
{
return;
}

_isScheduled = true;
}

_parent._scheduler.Schedule(this);
}

/// <summary>Stores the error for delivery after any pending value.</summary>
/// <param name="error">The error to store.</param>
public void OnError(Exception error)
{
lock (_gate)
{
if (_isDisposed || _error is not null)
{
return;
}

_error = error;
if (_isScheduled)
{
return;
}

_isScheduled = true;
}

_parent._scheduler.Schedule(this);
}

/// <summary>Stores completion for delivery after any pending value.</summary>
public void OnCompleted()
{
lock (_gate)
{
if (_isDisposed || _hasCompleted || _error is not null)
{
return;
}

_hasCompleted = true;
if (_isScheduled)
{
return;
}

_isScheduled = true;
}

_parent._scheduler.Schedule(this);
}

/// <summary>Executes the scheduled drain.</summary>
public void Execute()
{
while (true)
{
var value = default(T)!;
Exception? error = null;
var isValue = false;
lock (_gate)
{
if (_isDisposed)
{
_isScheduled = false;
return;
}

if (_hasPending)
{
value = _pending;
_pending = default!;
_hasPending = false;
isValue = true;
}
else if (_error is not null)
{
error = _error;
}
else if (!_hasCompleted)
{
_isScheduled = false;
return;
}
}

if (isValue)
{
_observer.OnNext(value);
continue;
}

if (error is not null)
{
_observer.OnError(error);
}
else
{
_observer.OnCompleted();
}

Dispose();
return;
}
}

/// <summary>Drops the pending value and releases the upstream subscription, once.</summary>
public void Dispose()
{
lock (_gate)
{
if (_isDisposed)
{
return;
}

_isDisposed = true;
_pending = default!;
_hasPending = false;
}

_cancel.Dispose();
}
}
}
33 changes: 33 additions & 0 deletions src/Primitives.Shared/Concurrency/SequencerImmediacyExtensions.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
// Copyright (c) 2019-2026 ReactiveUI Association Incorporated. All rights reserved.
// ReactiveUI Association Incorporated licenses this file to you under the MIT license.
// See the LICENSE file in the project root for full license information.

using System.Runtime.CompilerServices;

#if REACTIVE_SHIM
namespace ReactiveUI.Primitives.Reactive.Concurrency;
#else
namespace ReactiveUI.Primitives.Concurrency;
#endif

/// <summary>Reports whether a sequencer runs work inline on the calling thread.</summary>
public static class SequencerImmediacyExtensions
{
/// <summary>Immediacy check for a sequencer.</summary>
/// <param name="sequencer">The sequencer to test; may be <see langword="null"/>.</param>
extension(ISequencer? sequencer)
{
/// <summary>Gets a value indicating whether the sequencer is the immediate sequencer.</summary>
/// <value><see langword="true"/> when the sequencer runs work inline on the calling thread; <see langword="false"/> for any other sequencer or <see langword="null"/>.</value>
public bool IsImmediate
{
[MethodImpl(MethodImplOptions.AggressiveInlining)]
get =>
#if REACTIVE_SHIM
sequencer is global::System.Reactive.Concurrency.ImmediateScheduler;
#else
sequencer is ImmediateSequencer;
#endif
}
}
}
Loading
Loading