Skip to content

Commit a9fb343

Browse files
authored
fix(routing): deliver scheduled notifications in order so Navigate returns its view model (#4484)
## Summary **Awaiting `RoutingState.Navigate.Execute` returns the navigated view model on any main-thread scheduler.** - **Scheduled notifications arrive in order.** The internal observer that delivers results on a scheduler queues `OnNext`, `OnError` and `OnCompleted` and delivers them from a single scheduled drain, so a concurrent scheduler cannot run completion before the value. - **The same ordering applies to `ReactiveProperty` and `ScheduledSubject`,** which use the same observer. ## Why **With a concurrent main-thread scheduler, an awaited navigation could fail after it had succeeded.** - When no platform is registered, `RxSchedulers.MainThreadScheduler` is the task pool. Each notification was scheduled as its own work item, so `OnCompleted` could run first and the awaiter threw "The source completed without producing a value", although the view model was already on the navigation stack. Closes #4480 ## Breaking changes **None.** - `Sequencer.Immediate` still delivers inline. Other schedulers deliver a burst of notifications from one scheduled work item instead of one work item each. ## How this was verified **A test uses a scheduler that runs queued work newest first and checks that `Navigate` delivers the view model before it completes; a console app that awaits `Navigate.Execute` 1000 times on the task-pool scheduler no longer fails.** ## Notes for the reviewer **The fix is `SchedulingObserver`; the test is in `RoutingStateTests`.** - The queue is guarded by a lock, and a flag makes sure only one drain is scheduled at a time; the drain clears the flag when the queue is empty. - `ReactiveCommand` was not affected: `Execute()` hands its result to the caller on the execution thread. ## Checklist - [x] I have read the [Contribute guide](https://www.reactiveui.net/contribute/index.html) - [x] The PR title follows [Conventional Commits](https://www.conventionalcommits.org/) - [x] Tests cover this change, or the summary says why they do not - [x] New or changed public API has XML documentation
1 parent 37fa90a commit a9fb343

2 files changed

Lines changed: 236 additions & 24 deletions

File tree

‎src/ReactiveUI.Shared/Internal/Subjects/SchedulingObserver.cs‎

Lines changed: 115 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -11,42 +11,133 @@ namespace ReactiveUI.Reactive.Internal;
1111
#else
1212
namespace ReactiveUI.Internal;
1313
#endif
14-
/// Forwards each notification to a downstream observer on a scheduler, replacing a per-observer ObserveOn.
14+
///
15+
/// Forwards each notification to a downstream observer on a scheduler, replacing a per-observer ObserveOn.
16+
/// Notifications are queued and delivered in arrival order by a single scheduled drain, so a concurrent scheduler
17+
/// (such as the task pool) cannot run a later notification (for example OnCompleted) before an earlier one.
18+
///
1519
/// The element type.
1620
/// The observer that receives the scheduled notifications.
1721
/// The scheduler each notification is delivered on.
1822
internal sealed class SchedulingObserver<T>(IObserver<T> downstream, ISequencer scheduler) : IObserver<T>
1923
{
24+
/// Guards and .
25+
#if NET9_0_OR_GREATER
26+
private readonly Lock _gate = new();
27+
#else
28+
private readonly object _gate = new();
29+
#endif
30+
31+
/// Notifications waiting for the drain, in arrival order.
32+
private readonly Queue<Notification> _queue = new();
33+
34+
/// Whether a drain is scheduled or running.
35+
private bool _draining;
36+
37+
/// The kind of a queued notification.
38+
private enum NotificationKind
39+
{
40+
/// An OnNext value.
41+
Next = 0,
42+
43+
/// An OnError failure.
44+
Error = 1,
45+
46+
/// An OnCompleted signal.
47+
Completed = 2,
48+
}
49+
2050
///
2151
[MethodImpl(MethodImplOptions.AggressiveInlining)]
22-
public void OnNext(T value) =>
23-
scheduler.ScheduleOrInline(
24-
(Observer: downstream, Value: value),
25-
static (_, state) =>
26-
{
27-
state.Observer.OnNext(state.Value);
28-
return EmptyDisposable.Instance;
29-
});
52+
public void OnNext(T value) => Enqueue(new(NotificationKind.Next, value, null));
3053

3154
///
3255
[MethodImpl(MethodImplOptions.AggressiveInlining)]
33-
public void OnError(Exception error) =>
34-
scheduler.ScheduleOrInline(
35-
(Observer: downstream, Error: error),
36-
static (_, state) =>
37-
{
38-
state.Observer.OnError(state.Error);
39-
return EmptyDisposable.Instance;
40-
});
56+
public void OnError(Exception error) => Enqueue(new(NotificationKind.Error, default!, error));
4157

4258
///
4359
[MethodImpl(MethodImplOptions.AggressiveInlining)]
44-
public void OnCompleted() =>
45-
scheduler.ScheduleOrInline(
46-
downstream,
47-
static (_, observer) =>
60+
public void OnCompleted() => Enqueue(new(NotificationKind.Completed, default!, null));
61+
62+
/// Queues a notification and schedules a drain when none is pending.
63+
/// The notification to deliver.
64+
private void Enqueue(in Notification notification)
65+
{
66+
// The immediate scheduler runs work inline on the calling thread, so arrival order is already kept.
67+
if (ReferenceEquals(scheduler, Sequencer.Immediate))
68+
{
69+
Deliver(notification);
70+
return;
71+
}
72+
73+
lock (_gate)
74+
{
75+
_queue.Enqueue(notification);
76+
if (_draining)
77+
{
78+
return;
79+
}
80+
81+
_draining = true;
82+
}
83+
84+
_ = scheduler.Schedule(this, static (_, self) =>
85+
{
86+
self.Drain();
87+
return EmptyDisposable.Instance;
88+
});
89+
}
90+
91+
/// Delivers queued notifications in order until the queue is empty.
92+
private void Drain()
93+
{
94+
while (true)
95+
{
96+
Notification notification;
97+
lock (_gate)
4898
{
49-
observer.OnCompleted();
50-
return EmptyDisposable.Instance;
51-
});
99+
if (_queue.Count == 0)
100+
{
101+
_draining = false;
102+
return;
103+
}
104+
105+
notification = _queue.Dequeue();
106+
}
107+
108+
Deliver(notification);
109+
}
110+
}
111+
112+
/// Forwards one notification to the downstream observer.
113+
/// The notification to forward.
114+
private void Deliver(in Notification notification)
115+
{
116+
switch (notification.Kind)
117+
{
118+
case NotificationKind.Next:
119+
{
120+
downstream.OnNext(notification.Value);
121+
break;
122+
}
123+
124+
case NotificationKind.Error:
125+
{
126+
downstream.OnError(notification.Error!);
127+
break;
128+
}
129+
130+
default:
131+
{
132+
downstream.OnCompleted();
133+
break;
134+
}
135+
}
136+
}
137+
138+
/// A queued notification.
139+
/// The notification kind.
140+
/// The value, for .
141+
/// The failure, for .
142+
private readonly record struct Notification(NotificationKind Kind, T Value, Exception? Error);
52143
}
Lines changed: 121 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,121 @@
1+
// Copyright (c) 2009-2026 .NET Foundation and Contributors. All rights reserved.
2+
// Licensed to the .NET Foundation under one or more agreements.
3+
// The .NET Foundation licenses this file to you under the MIT license.
4+
// See the LICENSE file in the project root for full license information.
5+
6+
#if REACTIVE_SHIM
7+
using System.Reactive.Concurrency;
8+
#else
9+
using ReactiveUI.Primitives.Concurrency;
10+
#endif
11+
12+
namespace ReactiveUI.Tests.Routing;
13+
14+
/// Tests for .
15+
public class RoutingStateTests
16+
{
17+
///
18+
/// Navigate delivers the navigated view model before it completes, even when the navigation scheduler runs
19+
/// queued work out of order (as a concurrent scheduler such as the task pool can).
20+
///
21+
/// A representing the asynchronous operation.
22+
[Test]
23+
public async Task NavigateDeliversViewModelBeforeCompletionOnReorderingScheduler()
24+
{
25+
var scheduler = new ReorderingScheduler();
26+
var router = new RoutingState(scheduler);
27+
var viewModel = new TestViewModel();
28+
var notifications = new List<string>();
29+
30+
using var subscription = router.Navigate.Execute(viewModel).Subscribe(
31+
vm => notifications.Add(ReferenceEquals(vm, viewModel) ? "next" : "next-other"),
32+
() => notifications.Add("completed"));
33+
scheduler.RunNewestFirst();
34+
35+
await Assert.That(notifications).IsEquivalentTo(["next", "completed"], TUnit.Assertions.Enums.CollectionOrdering.Matching);
36+
}
37+
38+
/// A routable view model used as the navigation target.
39+
private sealed class TestViewModel : ReactiveObject, IRoutableViewModel
40+
{
41+
///
42+
public string? UrlPathSegment => "test";
43+
44+
///
45+
public IScreen HostScreen => null!;
46+
}
47+
48+
#if REACTIVE_SHIM
49+
///
50+
/// A scheduler that holds work until , then runs it newest first. It stands in for a
51+
/// concurrent scheduler, which gives no ordering guarantee between separately scheduled work items.
52+
///
53+
private sealed class ReorderingScheduler : IScheduler
54+
{
55+
/// The pending work, newest last.
56+
private readonly List<Action> _pending = [];
57+
58+
///
59+
public DateTimeOffset Now => default;
60+
61+
///
62+
public IDisposable Schedule<TState>(TState state, Func<IScheduler, TState, IDisposable> action)
63+
{
64+
_pending.Add(() => action(this, state));
65+
return EmptyDisposable.Instance;
66+
}
67+
68+
///
69+
public IDisposable Schedule<TState>(TState state, TimeSpan dueTime, Func<IScheduler, TState, IDisposable> action) =>
70+
Schedule(state, action);
71+
72+
///
73+
public IDisposable Schedule<TState>(TState state, DateTimeOffset dueTime, Func<IScheduler, TState, IDisposable> action) =>
74+
Schedule(state, action);
75+
76+
/// Runs pending work, newest first, until none remains.
77+
public void RunNewestFirst()
78+
{
79+
while (_pending.Count > 0)
80+
{
81+
var work = _pending[^1];
82+
_pending.RemoveAt(_pending.Count - 1);
83+
work();
84+
}
85+
}
86+
}
87+
#else
88+
///
89+
/// A scheduler that holds work until , then runs it newest first. It stands in for a
90+
/// concurrent scheduler, which gives no ordering guarantee between separately scheduled work items.
91+
///
92+
private sealed class ReorderingScheduler : ISequencer
93+
{
94+
/// The pending work, newest last.
95+
private readonly List<IWorkItem> _pending = [];
96+
97+
///
98+
public DateTimeOffset Now => default;
99+
100+
///
101+
public long Timestamp => 0;
102+
103+
///
104+
public void Schedule(IWorkItem item) => _pending.Add(item);
105+
106+
///
107+
public void Schedule(IWorkItem item, long dueTimestamp) => _pending.Add(item);
108+
109+
/// Runs pending work, newest first, until none remains.
110+
public void RunNewestFirst()
111+
{
112+
while (_pending.Count > 0)
113+
{
114+
var work = _pending[^1];
115+
_pending.RemoveAt(_pending.Count - 1);
116+
work.Execute();
117+
}
118+
}
119+
}
120+
#endif
121+
}

0 commit comments

Comments
 (0)