From 3c119339146caba2b1038d9de334433d7cd1a490 Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Wed, 26 Aug 2026 02:32:07 -0700 Subject: [PATCH 1/4] fix(test): decouple reminder barrier from scheduler --- .../ReminderService/LocalReminderService.cs | 59 +++++++++++-------- .../TimerTests/LocalReminderServiceTests.cs | 34 +++++++++++ 2 files changed, 69 insertions(+), 24 deletions(-) diff --git a/src/Orleans.Reminders/ReminderService/LocalReminderService.cs b/src/Orleans.Reminders/ReminderService/LocalReminderService.cs index 6b3f0b82e1d..6b290fa9037 100644 --- a/src/Orleans.Reminders/ReminderService/LocalReminderService.cs +++ b/src/Orleans.Reminders/ReminderService/LocalReminderService.cs @@ -31,6 +31,8 @@ internal sealed partial class LocalReminderService : GrainService, IReminderServ private readonly TimeProvider _timeProvider; private readonly ReminderInstruments _reminderInstruments; private long localTableSequence; + // The test barrier reads this state off-scheduler so it remains observable while the service is busy. + private readonly object _rangeChangeLock = new(); private long rangeChangeGeneration; private TaskCompletionSource rangeChangeGenerationChanged = new(TaskCreationOptions.RunContinuationsAsynchronously); private Task rangeChangeTask = Task.CompletedTask; @@ -397,18 +399,24 @@ public override Task OnRangeChange(IRingRange oldRange, IRingRange newRange, boo CheckRuntimeContext(); _ = base.OnRangeChange(oldRange, newRange, increased); - var task = Status == GrainServiceStatus.Started - ? ReadAndUpdateReminders() - : Task.CompletedTask; - if (Status != GrainServiceStatus.Started) - { - LogIgnoringRangeChange(Status); + Task task; + TaskCompletionSource previousGenerationChanged; + lock (_rangeChangeLock) + { + task = Status == GrainServiceStatus.Started + ? ReadAndUpdateReminders() + : Task.CompletedTask; + if (Status != GrainServiceStatus.Started) + { + LogIgnoringRangeChange(Status); + } + + previousGenerationChanged = rangeChangeGenerationChanged; + rangeChangeGenerationChanged = new(TaskCreationOptions.RunContinuationsAsynchronously); + rangeChangeGeneration++; + rangeChangeTask = task; } - var previousGenerationChanged = rangeChangeGenerationChanged; - rangeChangeGenerationChanged = new(TaskCreationOptions.RunContinuationsAsynchronously); - rangeChangeGeneration++; - rangeChangeTask = task; previousGenerationChanged.TrySetResult(); return task; } @@ -419,30 +427,33 @@ internal async Task TestOnlyWaitForRangeChangeReconciliation(CancellationToken c { // A newer refresh supersedes older results through localTableSequence. Follow generation // changes so a stalled obsolete read does not block the current reconciliation. - long observedGeneration = 0; - Task observedTask = Task.CompletedTask; - Task observedGenerationChanged = Task.CompletedTask; - await this.QueueTask(() => + long observedGeneration; + Task observedTask; + Task observedGenerationChanged; + lock (_rangeChangeLock) { observedGeneration = rangeChangeGeneration; observedTask = rangeChangeTask; observedGenerationChanged = rangeChangeGenerationChanged.Task; - return Task.CompletedTask; - }).WaitAsync(cancellationToken); + } await Task.WhenAny(observedTask, observedGenerationChanged).WaitAsync(cancellationToken); - var isCurrentGeneration = false; - await this.QueueTask(() => + lock (_rangeChangeLock) { - isCurrentGeneration = observedGeneration == rangeChangeGeneration; - return Task.CompletedTask; - }).WaitAsync(cancellationToken); + if (observedGeneration != rangeChangeGeneration) + { + continue; + } + } - if (isCurrentGeneration) + await observedTask.WaitAsync(cancellationToken); + lock (_rangeChangeLock) { - await observedTask.WaitAsync(cancellationToken); - return; + if (observedGeneration == rangeChangeGeneration) + { + return; + } } } } diff --git a/test/Orleans.Reminders.Tests/TimerTests/LocalReminderServiceTests.cs b/test/Orleans.Reminders.Tests/TimerTests/LocalReminderServiceTests.cs index a873b461e42..14f8886b2e1 100644 --- a/test/Orleans.Reminders.Tests/TimerTests/LocalReminderServiceTests.cs +++ b/test/Orleans.Reminders.Tests/TimerTests/LocalReminderServiceTests.cs @@ -328,6 +328,40 @@ public async Task RangeChangeBarrier_StopsWaitingWhenCanceled() } } + [TestSuite("BVT")] + [TestProvider("None")] + [Fact, TestCategory("BVT")] + public async Task RangeChangeBarrier_DoesNotDependOnServiceSchedulerAvailability() + { + var silo = Assert.Single(fixture.HostedCluster.Silos); + using var cancellation = new CancellationTokenSource(TestConstants.InitTimeout); + _ = await fixture.ReminderObserver.WaitForReminderServiceStartedAsync(cancellation.Token, silo.SiloAddress); + + var reminderService = silo.ServiceProvider.GetRequiredService(); + await reminderService.TestOnlyWaitForRangeChangeReconciliation(cancellation.Token); + + var schedulerBlocked = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + using var releaseScheduler = new ManualResetEventSlim(); + var blockingTask = new Task(() => + { + schedulerBlocked.TrySetResult(); + releaseScheduler.Wait(cancellation.Token); + }); + reminderService.Scheduler.QueueTask(blockingTask); + await schedulerBlocked.Task.WaitAsync(cancellation.Token); + + try + { + using var barrierCancellation = new CancellationTokenSource(TimeSpan.FromSeconds(1)); + await reminderService.TestOnlyWaitForRangeChangeReconciliation(barrierCancellation.Token); + } + finally + { + releaseScheduler.Set(); + await blockingTask.WaitAsync(cancellation.Token); + } + } + public sealed class Fixture : BaseInProcessTestClusterFixture { public ReminderDiagnosticObserver ReminderObserver { get; } = ReminderDiagnosticObserver.Create(); From 1995353edd8f6ecb238affb198a12961992e08c4 Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Wed, 26 Aug 2026 04:09:23 -0700 Subject: [PATCH 2/4] fix(reminders): narrow reconciliation barrier lock --- .../ReminderService/LocalReminderService.cs | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/src/Orleans.Reminders/ReminderService/LocalReminderService.cs b/src/Orleans.Reminders/ReminderService/LocalReminderService.cs index 6b290fa9037..5937ee5bde9 100644 --- a/src/Orleans.Reminders/ReminderService/LocalReminderService.cs +++ b/src/Orleans.Reminders/ReminderService/LocalReminderService.cs @@ -399,18 +399,18 @@ public override Task OnRangeChange(IRingRange oldRange, IRingRange newRange, boo CheckRuntimeContext(); _ = base.OnRangeChange(oldRange, newRange, increased); - Task task; + var status = Status; + var task = status == GrainServiceStatus.Started + ? ReadAndUpdateReminders() + : Task.CompletedTask; + if (status != GrainServiceStatus.Started) + { + LogIgnoringRangeChange(status); + } + TaskCompletionSource previousGenerationChanged; lock (_rangeChangeLock) { - task = Status == GrainServiceStatus.Started - ? ReadAndUpdateReminders() - : Task.CompletedTask; - if (Status != GrainServiceStatus.Started) - { - LogIgnoringRangeChange(Status); - } - previousGenerationChanged = rangeChangeGenerationChanged; rangeChangeGenerationChanged = new(TaskCreationOptions.RunContinuationsAsynchronously); rangeChangeGeneration++; From 2a4d9e99a39b94d7b6cf4805f9b7accc00a71847 Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Wed, 26 Aug 2026 07:32:47 -0700 Subject: [PATCH 3/4] fix(reminders): publish reconciliation before reads --- .../ReminderService/LocalReminderService.cs | 32 ++++++++++++------- 1 file changed, 21 insertions(+), 11 deletions(-) diff --git a/src/Orleans.Reminders/ReminderService/LocalReminderService.cs b/src/Orleans.Reminders/ReminderService/LocalReminderService.cs index 5937ee5bde9..6088a4a2c7b 100644 --- a/src/Orleans.Reminders/ReminderService/LocalReminderService.cs +++ b/src/Orleans.Reminders/ReminderService/LocalReminderService.cs @@ -399,26 +399,36 @@ public override Task OnRangeChange(IRingRange oldRange, IRingRange newRange, boo CheckRuntimeContext(); _ = base.OnRangeChange(oldRange, newRange, increased); - var status = Status; - var task = status == GrainServiceStatus.Started - ? ReadAndUpdateReminders() - : Task.CompletedTask; - if (status != GrainServiceStatus.Started) - { - LogIgnoringRangeChange(status); - } - + var reconciliationTaskSource = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); TaskCompletionSource previousGenerationChanged; lock (_rangeChangeLock) { previousGenerationChanged = rangeChangeGenerationChanged; rangeChangeGenerationChanged = new(TaskCreationOptions.RunContinuationsAsynchronously); rangeChangeGeneration++; - rangeChangeTask = task; + rangeChangeTask = reconciliationTaskSource.Task.Unwrap(); } previousGenerationChanged.TrySetResult(); - return task; + try + { + var status = Status; + var task = status == GrainServiceStatus.Started + ? ReadAndUpdateReminders() + : Task.CompletedTask; + if (status != GrainServiceStatus.Started) + { + LogIgnoringRangeChange(status); + } + + reconciliationTaskSource.SetResult(task); + return task; + } + catch (Exception exception) + { + reconciliationTaskSource.SetException(exception); + throw; + } } internal async Task TestOnlyWaitForRangeChangeReconciliation(CancellationToken cancellationToken) From d99913f4b123a27c73a656a11c87549a82f4b47c Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Wed, 26 Aug 2026 08:23:19 -0700 Subject: [PATCH 4/4] fix(reminders): ignore obsolete reconciliation faults --- .../ReminderService/LocalReminderService.cs | 12 ++++-------- 1 file changed, 4 insertions(+), 8 deletions(-) diff --git a/src/Orleans.Reminders/ReminderService/LocalReminderService.cs b/src/Orleans.Reminders/ReminderService/LocalReminderService.cs index 6088a4a2c7b..d66c96592ae 100644 --- a/src/Orleans.Reminders/ReminderService/LocalReminderService.cs +++ b/src/Orleans.Reminders/ReminderService/LocalReminderService.cs @@ -455,15 +455,11 @@ internal async Task TestOnlyWaitForRangeChangeReconciliation(CancellationToken c { continue; } - } - await observedTask.WaitAsync(cancellationToken); - lock (_rangeChangeLock) - { - if (observedGeneration == rangeChangeGeneration) - { - return; - } + // The generation-change task can only complete after the generation advances, so the + // observed reconciliation is complete and its outcome can be read without blocking. + observedTask.GetAwaiter().GetResult(); + return; } } }