From 10c6c045d783b6daf434eca9617b4493f731707b Mon Sep 17 00:00:00 2001
From: Glenn Watson <5834289+glennawatson@users.noreply.github.com>
Date: Wed, 2 Sep 2026 05:20:33 +1000
Subject: [PATCH 1/2] test: drive witness and from-task dispatch from gates
instead of the clock
- WitnessOn is not thread-pool specialized, so its test now drives a
ManualSequencer and asserts dispatch is deferred until the sequencer drains
it, replacing a real pool subscription behind a 2s wall clock and a
`values.Count <= 1` assertion that held even when every value was dropped.
- The from-task completion and user-exception tests raced a 10s in-command
delay against a 10.5s/11.5s test-side delay. They now release the command
body through a TaskCompletionSource and wait on the terminal cleanup, the
pattern the file's other from-task tests already use.
- Dispose in the user-exception test now lands after the terminal has been
observed, which is the ordering its assertions always described.
---
.../SignalFromTaskTest.cs | 72 ++++++++++++++-----
.../WitnessTests.cs | 54 +++++++++-----
2 files changed, 91 insertions(+), 35 deletions(-)
diff --git a/src/tests/ReactiveUI.Primitives.Tests/SignalFromTaskTest.cs b/src/tests/ReactiveUI.Primitives.Tests/SignalFromTaskTest.cs
index 2a2eb73d..2a6e8383 100644
--- a/src/tests/ReactiveUI.Primitives.Tests/SignalFromTaskTest.cs
+++ b/src/tests/ReactiveUI.Primitives.Tests/SignalFromTaskTest.cs
@@ -50,9 +50,6 @@ public class SignalFromTaskTest
/// Delay used by the command body.
private const int CommandDelayMilliseconds = 10_000;
- /// Delay used to wait for normal command completion.
- private const int CompletionWaitDelayMilliseconds = 11_000;
-
/// Exception message used by user exception tests.
private const string BreakExecutionMessage = "break execution";
@@ -513,17 +510,25 @@ public void ImmediateSignalSubscribeAfterDisposeThrows()
_ = Assert.Throws(() => taskSignal.Subscribe(static _ => { }));
}
- /// Signals from task handles user exceptions.
+ ///
+ /// Signals from task handles user exceptions. The command body is released by a gate rather than a timer, so the
+ /// subscription is only torn down once the failure has already travelled the whole chain; a dispose that lands
+ /// after the terminal must leave the cancellation path untouched.
+ ///
/// A representing the asynchronous unit test.
[Test]
public async Task SignalFromTaskHandlesUserExceptions()
{
StatusTrail statusTrail = new();
+ TaskCompletionSource executionStarted = new(TaskCreationOptions.RunContinuationsAsynchronously);
+ TaskCompletionSource releaseExecution = new(TaskCreationOptions.RunContinuationsAsynchronously);
+ TaskCompletionSource finallyCompleted = new(TaskCreationOptions.RunContinuationsAsynchronously);
var position = 0;
var fixture = Signal.FromTask(async cts =>
{
RecordStatus(statusTrail, ref position, StartedCommand);
- await Task.Delay(CommandDelayMilliseconds, cts.Token)
+ _ = executionStarted.TrySetResult();
+ await releaseExecution.Task.WaitAsync(cts.Token)
.HandleCancellation(() => RecordCancellationCleanup(statusTrail, ref position)).ConfigureAwait(true);
if (!cts.IsCancellationRequested)
{
@@ -535,14 +540,18 @@ await Task.Delay(CommandDelayMilliseconds, cts.Token)
{
RecordStatus(statusTrail, ref position, ExceptionShouldBeHere);
return Signal.Fail(ex);
- }).OnCleanup(() => RecordStatus(statusTrail, ref position, ShouldAlwaysComeHere));
+ }).OnCleanup(() =>
+ {
+ RecordStatus(statusTrail, ref position, ShouldAlwaysComeHere);
+ _ = finallyCompleted.TrySetResult();
+ });
var result = false;
var subscription = fixture.Subscribe(_ => result = true);
- await Task.Delay(InitialDelayMilliseconds).ConfigureAwait(true);
+ await executionStarted.Task.WaitAsync(PollTimeout).ConfigureAwait(false);
await Assert.That(StatusMessages(statusTrail)).Contains(StartedCommand);
- await Task.Delay(CommandDelayMilliseconds).ConfigureAwait(true);
+ _ = releaseExecution.TrySetResult();
+ await finallyCompleted.Task.WaitAsync(PollTimeout).ConfigureAwait(false);
subscription.Dispose();
- await Task.Delay(CancellationWaitDelayMilliseconds).ConfigureAwait(false);
await Assert.That(StatusMessages(statusTrail)).DoesNotContain(StartingCancellingCommand);
await Assert.That(StatusMessages(statusTrail)).Contains(ShouldAlwaysComeHere);
await Assert.That(StatusMessages(statusTrail)).DoesNotContain(FinishedCancellingCommand);
@@ -666,17 +675,25 @@ public async Task SignalFromTaskHandlesCancellationInBase()
await Assert.That(statusTrail.LastMessage).IsEqualTo(ShouldAlwaysComeHere);
}
- /// Signals from task handles completion.
+ ///
+ /// Signals from task handles completion. The command body is released by a gate rather than a timer, so the
+ /// assertions run once the terminal cleanup has actually happened instead of once a wall-clock window is judged
+ /// long enough for it.
+ ///
/// A representing the asynchronous unit test.
[Test]
public async Task SignalFromTaskHandlesCompletion()
{
StatusTrail statusTrail = new();
+ TaskCompletionSource executionStarted = new(TaskCreationOptions.RunContinuationsAsynchronously);
+ TaskCompletionSource releaseExecution = new(TaskCreationOptions.RunContinuationsAsynchronously);
+ TaskCompletionSource finallyCompleted = new(TaskCreationOptions.RunContinuationsAsynchronously);
var position = 0;
var fixture = Signal.FromTask(async cts =>
{
RecordStatus(statusTrail, ref position, StartedCommand);
- await Task.Delay(CommandDelayMilliseconds, cts.Token)
+ _ = executionStarted.TrySetResult();
+ await releaseExecution.Task.WaitAsync(cts.Token)
.HandleCancellation(() => RecordCancellationCleanup(statusTrail, ref position)).ConfigureAwait(true);
if (!cts.IsCancellationRequested)
{
@@ -688,12 +705,17 @@ await Task.Delay(CommandDelayMilliseconds, cts.Token)
{
RecordStatus(statusTrail, ref position, ExceptionShouldBeHere);
return Signal.Fail(ex);
- }).OnCleanup(() => RecordStatus(statusTrail, ref position, ShouldAlwaysComeHere));
+ }).OnCleanup(() =>
+ {
+ RecordStatus(statusTrail, ref position, ShouldAlwaysComeHere);
+ _ = finallyCompleted.TrySetResult();
+ });
var result = false;
using var subscription = fixture.Subscribe(_ => result = true);
- await Task.Delay(InitialDelayMilliseconds).ConfigureAwait(true);
+ await executionStarted.Task.WaitAsync(PollTimeout).ConfigureAwait(false);
await Assert.That(StatusMessages(statusTrail)).Contains(StartedCommand);
- await Task.Delay(CompletionWaitDelayMilliseconds).ConfigureAwait(false);
+ _ = releaseExecution.TrySetResult();
+ await finallyCompleted.Task.WaitAsync(PollTimeout).ConfigureAwait(false);
await Assert.That(StatusMessages(statusTrail)).DoesNotContain(StartingCancellingCommand);
await Assert.That(StatusMessages(statusTrail)).DoesNotContain(FinishedCancellingCommand);
await Assert.That(StatusMessages(statusTrail)).Contains(FinishedCommandNormally);
@@ -861,17 +883,24 @@ public async Task SignalFromTask_T_HandlesCancellationInBase()
await Assert.That(statusTrail.LastMessage).IsEqualTo(ShouldAlwaysComeHere);
}
- /// Signals from task t handles completion.
+ ///
+ /// Signals from task t handles completion. Like its non-generic counterpart the command body is released by a
+ /// gate rather than a timer, so the assertions run once the terminal cleanup has actually happened.
+ ///
/// A representing the asynchronous unit test.
[Test]
public async Task SignalFromTask_T_HandlesCompletion()
{
StatusTrail statusTrail = new();
+ TaskCompletionSource executionStarted = new(TaskCreationOptions.RunContinuationsAsynchronously);
+ TaskCompletionSource releaseExecution = new(TaskCreationOptions.RunContinuationsAsynchronously);
+ TaskCompletionSource finallyCompleted = new(TaskCreationOptions.RunContinuationsAsynchronously);
var position = 0;
var fixture = Signal.FromTask(async cts =>
{
RecordStatus(statusTrail, ref position, StartedCommand);
- await Task.Delay(CommandDelayMilliseconds, cts.Token)
+ _ = executionStarted.TrySetResult();
+ await releaseExecution.Task.WaitAsync(cts.Token)
.HandleCancellation(() => RecordCancellationCleanup(statusTrail, ref position)).ConfigureAwait(true);
if (!cts.IsCancellationRequested)
{
@@ -883,12 +912,17 @@ await Task.Delay(CommandDelayMilliseconds, cts.Token)
{
RecordStatus(statusTrail, ref position, ExceptionShouldBeHere);
return Signal.Fail(ex);
- }).OnCleanup(() => RecordStatus(statusTrail, ref position, ShouldAlwaysComeHere));
+ }).OnCleanup(() =>
+ {
+ RecordStatus(statusTrail, ref position, ShouldAlwaysComeHere);
+ _ = finallyCompleted.TrySetResult();
+ });
var result = false;
using var subscription = fixture.Subscribe(_ => result = true);
- await Task.Delay(InitialDelayMilliseconds).ConfigureAwait(true);
+ await executionStarted.Task.WaitAsync(PollTimeout).ConfigureAwait(false);
await Assert.That(StatusMessages(statusTrail)).Contains(StartedCommand);
- await Task.Delay(CompletionWaitDelayMilliseconds).ConfigureAwait(false);
+ _ = releaseExecution.TrySetResult();
+ await finallyCompleted.Task.WaitAsync(PollTimeout).ConfigureAwait(false);
await Assert.That(StatusMessages(statusTrail)).DoesNotContain(StartingCancellingCommand);
await Assert.That(StatusMessages(statusTrail)).DoesNotContain(FinishedCancellingCommand);
await Assert.That(StatusMessages(statusTrail)).Contains(FinishedCommandNormally);
diff --git a/src/tests/ReactiveUI.Primitives.Tests/WitnessTests.cs b/src/tests/ReactiveUI.Primitives.Tests/WitnessTests.cs
index 0c0c7c96..669c4c8e 100644
--- a/src/tests/ReactiveUI.Primitives.Tests/WitnessTests.cs
+++ b/src/tests/ReactiveUI.Primitives.Tests/WitnessTests.cs
@@ -37,7 +37,7 @@ public class WitnessTests
/// A reusable value for fourteen.
private const int Fourteen = 14;
- /// Timeout used when waiting for thread-pool scheduled observer callbacks.
+ /// Timeout used when awaiting a witness task that has already been driven to its terminal.
private const int TimeoutSeconds = 2;
/// Shared state value.
@@ -52,7 +52,7 @@ public class WitnessTests
/// Expected safe witness event sequence.
private static readonly string[] ExpectedSafeEvents = ["next:3", "completed"];
- /// Expected values from thread-pool observer dispatch.
+ /// Expected values from sequenced observer dispatch.
private static readonly int[] WitnessOnExpected = [One];
/// Verifies delegate witnesses route next, error, and completion callbacks.
@@ -157,28 +157,50 @@ public async Task WitnessesCoverDisposedThrowEmptyAndSafeBranches()
_ = Assert.Throws(() => safe.OnError(null!));
}
- /// Covers the thread-pool-specialized witness dispatch implementation.
- /// A task representing asynchronous observer dispatch.
+ ///
+ /// Verifies the witness holds every notification back until its sequencer runs the queued drain, then replays
+ /// values and the completion through the observer in order. The dispatch is driven by a sequencer the test owns
+ /// so the handover is observed exactly rather than raced against a pool thread.
+ ///
+ /// A task representing the asynchronous operation.
[Test]
- public async Task WitnessOnThreadPoolDispatchesNextCompletedAndErrorSignals()
+ public async Task WitnessOnDefersNextAndCompletedUntilTheSequencerDrainsThem()
{
+ ManualSequencer sequencer = new();
List values = [];
- TaskCompletionSource completion = new(TaskCreationOptions.RunContinuationsAsynchronously);
+ var completed = 0;
+
using (Signal.FromEnumerable(WitnessOnExpected)
- .WitnessOn(ThreadPoolSequencer.Instance)
- .Subscribe(values.Add, completion.SetException, completion.SetResult))
+ .WitnessOn(sequencer)
+ .Subscribe(values.Add, static error => throw error, () => completed++))
{
- await WaitForAsync(completion.Task);
+ await Assert.That(values).IsEmpty();
+ await Assert.That(completed).IsEqualTo(0);
+ sequencer.RunPending();
}
- await Assert.That(values.Count <= WitnessOnExpected.Length).IsTrue();
- InvalidOperationException error = new("thread-pool");
- TaskCompletionSource observed = new(TaskCreationOptions.RunContinuationsAsynchronously);
- using (Signal.Fail(error).WitnessOn(ThreadPoolSequencer.Instance)
- .Subscribe(static _ => { }, observed.SetResult, static () => { }))
+ await Assert.That(values.SequenceEqual(WitnessOnExpected)).IsTrue();
+ await Assert.That(completed).IsEqualTo(1);
+ }
+
+ /// Verifies the witness routes a source failure through the same deferred sequencer drain.
+ /// A task representing the asynchronous operation.
+ [Test]
+ public async Task WitnessOnDefersAnErrorUntilTheSequencerDrainsIt()
+ {
+ ManualSequencer sequencer = new();
+ InvalidOperationException error = new("sequenced");
+ Exception? observed = null;
+
+ using (Signal.Fail(error)
+ .WitnessOn(sequencer)
+ .Subscribe(static _ => { }, failure => observed = failure, static () => { }))
{
- await Assert.That(await WaitForAsync(observed.Task)).IsSameReferenceAs(error);
+ await Assert.That(observed).IsNull();
+ sequencer.RunPending();
}
+
+ await Assert.That(observed).IsSameReferenceAs(error);
}
/// Covers callback, forwarding, and stateful witness contracts.
@@ -1044,7 +1066,7 @@ private static async Task WaitForAsync(Task task)
var completed = await Task.WhenAny(task, timeout).ConfigureAwait(false);
if (completed == timeout)
{
- throw new TimeoutException("Timed out waiting for scheduled observer dispatch.");
+ throw new TimeoutException("Timed out waiting for the witness task to complete.");
}
await task.ConfigureAwait(false);
From fa4f422d1bb5f6fa4163439803c5f7778c5f1fc4 Mon Sep 17 00:00:00 2001
From: Glenn Watson <5834289+glennawatson@users.noreply.github.com>
Date: Wed, 2 Sep 2026 06:14:29 +1000
Subject: [PATCH 2/2] fix: stop the work item and expire tests depending on a
free pool thread
- Extract the race-only tail of DelegateWorkItem.Execute into an excluded
helper. Single-threaded it is unreachable: a completed Dispose leaves the slot
holding EmptyDisposable, so the compare-exchange takes the already-claimed
branch and never falls through. It was only ever reached when the existing
Barrier test happened to win its race, which made the line's coverage a
property of machine load.
- ExpireCoordinatorTests dispatched a deliberately blocking OnNext through the
thread pool, so a saturated pool left the notification queued behind unrelated
work and the test timed out waiting to enter the observer. Both sides of the
race now run on dedicated threads.
- The same test read the invariant flag while both threads were still running,
through plain fields written on other threads. Those are volatile now, so the
assertion observes the write rather than relying on one being published.
---
.../Concurrency/Sequencer.Simple.cs | 16 +++++
.../ExpireCoordinatorTests.cs | 72 ++++++++++++++-----
2 files changed, 72 insertions(+), 16 deletions(-)
diff --git a/src/ReactiveUI.Primitives/Concurrency/Sequencer.Simple.cs b/src/ReactiveUI.Primitives/Concurrency/Sequencer.Simple.cs
index 8523b85c..ee9d2696 100644
--- a/src/ReactiveUI.Primitives/Concurrency/Sequencer.Simple.cs
+++ b/src/ReactiveUI.Primitives/Concurrency/Sequencer.Simple.cs
@@ -2,6 +2,7 @@
// 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.Diagnostics.CodeAnalysis;
using System.Runtime.CompilerServices;
using ReactiveUI.Primitives.Disposables;
@@ -140,6 +141,21 @@ public void Execute()
return;
}
+ DisposeIfRaced(disposable);
+ }
+
+ ///
+ /// Race-only cleanup: releases what the action returned when latched the
+ /// flag after the store above claimed the slot. Single-threaded this can never fire - a completed
+ /// leaves the slot holding , so the
+ /// compare-exchange takes the already-claimed branch instead and never reaches here. Only a real
+ /// concurrent disposal lands in this window, so it is excluded rather than chased with a
+ /// timing-dependent test.
+ ///
+ /// The disposable the scheduled action returned.
+ [ExcludeFromCodeCoverage]
+ private void DisposeIfRaced(IDisposable disposable)
+ {
if (!IsDisposed)
{
return;
diff --git a/src/tests/ReactiveUI.Primitives.Tests/ExpireCoordinatorTests.cs b/src/tests/ReactiveUI.Primitives.Tests/ExpireCoordinatorTests.cs
index 4a2c44ce..f0a4d714 100644
--- a/src/tests/ReactiveUI.Primitives.Tests/ExpireCoordinatorTests.cs
+++ b/src/tests/ReactiveUI.Primitives.Tests/ExpireCoordinatorTests.cs
@@ -179,17 +179,20 @@ public async Task TimeoutDoesNotEnterObserverWhileOnNextIsInFlight()
BlockingObserver observer = new();
using var subscription = source.Expire(TimeSpan.FromTicks(One), clock).Subscribe(observer);
- var onNextTask = Task.Run(() => source.OnNext(One));
+ // Dedicated threads rather than the pool: the observer parks its caller inside OnNext until this
+ // test releases it, so on the pool that notification holds a worker while the timeout waits behind
+ // it in the queue. A saturated pool then starves the very interleaving under test.
+ var onNextFinished = RunOnDedicatedThread(() => source.OnNext(One));
await observer.OnNextEntered.Task.WaitAsync(WaitTimeout).ConfigureAwait(false);
- var timeoutTask = Task.Run(() => clock.AdvanceBy(TimeSpan.FromTicks(One)));
+ var timeoutFinished = RunOnDedicatedThread(() => clock.AdvanceBy(TimeSpan.FromTicks(One)));
await Task.Delay(RaceSettleDelay).ConfigureAwait(false);
await Assert.That(observer.ErrorEnteredDuringOnNext).IsFalse();
observer.ReleaseOnNext.Set();
- await onNextTask.WaitAsync(WaitTimeout).ConfigureAwait(false);
- await timeoutTask.WaitAsync(WaitTimeout).ConfigureAwait(false);
+ await onNextFinished.WaitAsync(WaitTimeout).ConfigureAwait(false);
+ await timeoutFinished.WaitAsync(WaitTimeout).ConfigureAwait(false);
// Timeout may be observed after OnNext exits depending on scheduler timing.
// The invariant required here is that OnError never re-enters while OnNext is active.
@@ -197,6 +200,32 @@ public async Task TimeoutDoesNotEnterObserverWhileOnNextIsInFlight()
await Assert.That(observer.Values).IsEqualTo(One);
}
+ ///
+ /// Runs work on its own thread and reports when it finished. Used where the work blocks for the duration of
+ /// the scenario, which the thread pool cannot absorb without the risk of the work never being given a thread.
+ ///
+ /// The work to run.
+ /// A task that completes when the work has returned.
+ private static Task RunOnDedicatedThread(Action work)
+ {
+ TaskCompletionSource finished = new(TaskCreationOptions.RunContinuationsAsynchronously);
+ Thread thread = new(() =>
+ {
+ try
+ {
+ work();
+ _ = finished.TrySetResult();
+ }
+ catch (Exception error)
+ {
+ _ = finished.TrySetException(error);
+ }
+ }) { IsBackground = true };
+
+ thread.Start();
+ return finished.Task;
+ }
+
///
/// A sequencer that accepts scheduled work and never dispatches it, modelling a thread-pool sequencer whose pool
/// is saturated: the timer becomes due on the clock, but no thread is free to run the callback. Its clock is
@@ -245,6 +274,20 @@ private sealed class UndispatchedSequencer(DateTimeOffset start) : ISequencer
+ "not IDisposable.")]
private sealed class BlockingObserver : IObserver
{
+ /// Non-zero while is active. Written by the notifying thread and read by
+ /// the timeout thread, so the two must not race on a plain field.
+ private int _isInOnNext;
+
+ /// Non-zero once an error arrived while was active. The test reads this
+ /// while both threads are still running, so the write has to be published rather than merely made.
+ private int _errorEnteredDuringOnNext;
+
+ /// The number of forwarded values.
+ private int _values;
+
+ /// The number of forwarded errors.
+ private int _errors;
+
/// Gets the task completed when is entered.
public TaskCompletionSource OnNextEntered { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously);
@@ -252,16 +295,13 @@ private sealed class BlockingObserver : IObserver
public ManualResetEventSlim ReleaseOnNext { get; } = new();
/// Gets the number of forwarded values.
- public int Values { get; private set; }
+ public int Values => Volatile.Read(ref _values);
/// Gets the number of forwarded errors.
- public int Errors { get; private set; }
+ public int Errors => Volatile.Read(ref _errors);
/// Gets a value indicating whether an error entered while was active.
- public bool ErrorEnteredDuringOnNext { get; private set; }
-
- /// Gets or sets a value indicating whether is active.
- private bool IsInOnNext { get; set; }
+ public bool ErrorEnteredDuringOnNext => Volatile.Read(ref _errorEnteredDuringOnNext) != 0;
///
public void OnCompleted()
@@ -271,22 +311,22 @@ public void OnCompleted()
///
public void OnError(Exception error)
{
- if (IsInOnNext)
+ if (Volatile.Read(ref _isInOnNext) != 0)
{
- ErrorEnteredDuringOnNext = true;
+ Volatile.Write(ref _errorEnteredDuringOnNext, 1);
}
- Errors++;
+ _ = Interlocked.Increment(ref _errors);
}
///
public void OnNext(int value)
{
- Values++;
- IsInOnNext = true;
+ _ = Interlocked.Increment(ref _values);
+ Volatile.Write(ref _isInOnNext, 1);
OnNextEntered.SetResult();
_ = ReleaseOnNext.Wait(WaitTimeout);
- IsInOnNext = false;
+ Volatile.Write(ref _isInOnNext, 0);
}
}
}