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
67 changes: 42 additions & 25 deletions src/Orleans.Reminders/ReminderService/LocalReminderService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<Task>(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)
Expand All @@ -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;
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<LocalReminderService>();
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();
Expand Down