diff --git a/src/Orleans.Reminders/ReminderService/LocalReminderService.cs b/src/Orleans.Reminders/ReminderService/LocalReminderService.cs index 6b3f0b82e1d..d66c96592ae 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,20 +399,36 @@ 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); + var reconciliationTaskSource = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + TaskCompletionSource previousGenerationChanged; + lock (_rangeChangeLock) + { + previousGenerationChanged = rangeChangeGenerationChanged; + rangeChangeGenerationChanged = new(TaskCreationOptions.RunContinuationsAsynchronously); + rangeChangeGeneration++; + rangeChangeTask = reconciliationTaskSource.Task.Unwrap(); } - var previousGenerationChanged = rangeChangeGenerationChanged; - rangeChangeGenerationChanged = new(TaskCreationOptions.RunContinuationsAsynchronously); - rangeChangeGeneration++; - rangeChangeTask = task; 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) @@ -419,29 +437,28 @@ 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); + // 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; } } 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();