Skip to content

Commit cb1a19d

Browse files
committed
fix(testing): await reminder ring listener updates
1 parent d27da7c commit cb1a19d

4 files changed

Lines changed: 127 additions & 1 deletion

File tree

src/Orleans.Reminders/ReminderService/LocalReminderService.cs

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@
1010
using Orleans.Reminders.Diagnostics;
1111
using Orleans.Runtime.ConsistentRing;
1212
using Orleans.Runtime.Internal;
13+
using Orleans.Runtime.MembershipService;
1314
using Orleans.Runtime.Scheduler;
1415

1516
namespace Orleans.Runtime.ReminderService
@@ -28,6 +29,7 @@ internal sealed partial class LocalReminderService : GrainService, IReminderServ
2829
private readonly IAsyncTimer listRefreshTimer; // timer that refreshes our list of reminders to reflect global reminder table
2930
private readonly GrainReferenceActivator _referenceActivator;
3031
private readonly GrainInterfaceType _grainInterfaceType;
32+
private readonly SiloStatusListenerManager _siloStatusListenerManager;
3133
private readonly TimeProvider _timeProvider;
3234
private readonly ReminderInstruments _reminderInstruments;
3335
private long localTableSequence;
@@ -50,6 +52,7 @@ public LocalReminderService(
5052
IAsyncTimerFactory asyncTimerFactory,
5153
IOptions<ReminderOptions> reminderOptions,
5254
IConsistentRingProvider ringProvider,
55+
SiloStatusListenerManager siloStatusListenerManager,
5356
[FromKeyedServices(ReminderTimeProviderNames.Reminders)] TimeProvider timeProvider,
5457
ReminderInstruments reminderInstruments,
5558
SystemTargetShared shared)
@@ -60,6 +63,7 @@ public LocalReminderService(
6063
{
6164
_referenceActivator = referenceActivator;
6265
_grainInterfaceType = interfaceTypeResolver.GetGrainInterfaceType(typeof(IRemindable));
66+
_siloStatusListenerManager = siloStatusListenerManager;
6367
this.reminderOptions = reminderOptions.Value;
6468
this.reminderTable = reminderTable;
6569
_timeProvider = timeProvider;
@@ -370,6 +374,9 @@ internal Task TestOnlyRefresh()
370374
return refreshTask.Unwrap();
371375
}
372376

377+
internal Task TestOnlyWaitForSiloStatusListeners(CancellationToken cancellationToken)
378+
=> _siloStatusListenerManager.TestOnlyWaitForCurrentMembershipVersion(cancellationToken);
379+
373380
private void RemoveOutOfRangeReminders(List<Task> removedReminderTasks)
374381
{
375382
CheckRuntimeContext();

src/Orleans.Runtime/MembershipService/SiloStatusListenerManager.cs

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,10 @@ internal sealed partial class SiloStatusListenerManager : ILifecycleParticipant<
2222
private readonly IMembershipManager _membershipService;
2323
private readonly ILogger<SiloStatusListenerManager> _logger;
2424
private readonly IFatalErrorHandler _fatalErrorHandler;
25+
private readonly object _processedMembershipLock = new();
2526
private ImmutableList<WeakReference<ISiloStatusListener>> _listeners = [];
27+
private MembershipVersion _processedMembershipVersion = MembershipVersion.MinValue;
28+
private TaskCompletionSource _processedMembershipVersionChanged = CreateCompletion();
2629

2730
public SiloStatusListenerManager(
2831
IMembershipManager membershipManager,
@@ -88,6 +91,7 @@ private async Task ProcessMembershipUpdates()
8891
var update = (previous is null || snapshot.Version == MembershipVersion.MinValue) ? snapshot.AsUpdate() : snapshot.CreateUpdate(previous);
8992
NotifyObservers(update);
9093
previous = snapshot;
94+
MarkMembershipVersionProcessed(snapshot.Version);
9195
}
9296
}
9397
catch (OperationCanceledException) when (_cancellation.IsCancellationRequested)
@@ -144,6 +148,47 @@ private void NotifyObservers(ClusterMembershipUpdate update)
144148
}
145149
}
146150

151+
internal async Task TestOnlyWaitForCurrentMembershipVersion(CancellationToken cancellationToken)
152+
{
153+
var targetVersion = _membershipService.CurrentSnapshot.Version;
154+
while (true)
155+
{
156+
Task versionChanged;
157+
lock (_processedMembershipLock)
158+
{
159+
if (_processedMembershipVersion >= targetVersion)
160+
{
161+
return;
162+
}
163+
164+
versionChanged = _processedMembershipVersionChanged.Task;
165+
}
166+
167+
await versionChanged.WaitAsync(cancellationToken);
168+
}
169+
}
170+
171+
private void MarkMembershipVersionProcessed(MembershipVersion version)
172+
{
173+
TaskCompletionSource versionChanged;
174+
lock (_processedMembershipLock)
175+
{
176+
if (version <= _processedMembershipVersion)
177+
{
178+
return;
179+
}
180+
181+
_processedMembershipVersion = version;
182+
versionChanged = _processedMembershipVersionChanged;
183+
_processedMembershipVersionChanged = CreateCompletion();
184+
}
185+
186+
versionChanged.TrySetResult();
187+
}
188+
189+
private static TaskCompletionSource CreateCompletion()
190+
=> new(TaskCreationOptions.RunContinuationsAsynchronously);
191+
147192
void ILifecycleParticipant<ISiloLifecycle>.Participate(ISiloLifecycle lifecycle)
148193
{
149194
Task? task = null;

test/Orleans.Reminders.Tests/TimerTests/LocalReminderServiceTests.cs

Lines changed: 70 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -254,6 +254,51 @@ public async Task InitialRead_TreatsNullTableResultAsNoWork()
254254
Assert.True(reminderTable.RangeReadCount > 0);
255255
}
256256

257+
[TestSuite("BVT")]
258+
[TestProvider("None")]
259+
[Fact, TestCategory("BVT")]
260+
public async Task SiloStatusListenerBarrier_WaitsForCurrentMembershipVersion()
261+
{
262+
var initialSilo = Assert.Single(fixture.HostedCluster.Silos);
263+
using var cancellation = CancellationTokenSource.CreateLinkedTokenSource(TestContext.Current.CancellationToken);
264+
cancellation.CancelAfter(TestConstants.InitTimeout);
265+
var listener = new BlockingSiloStatusListener(initialSilo.SiloAddress);
266+
var statusOracle = initialSilo.ServiceProvider.GetRequiredService<ISiloStatusOracle>();
267+
Assert.True(statusOracle.SubscribeToSiloStatusEvents(listener));
268+
Task<List<InProcessSiloHandle>>? startTask = null;
269+
InProcessSiloHandle? joinedSilo = null;
270+
try
271+
{
272+
startTask = fixture.HostedCluster.StartSilosAsync(1, cancellation.Token);
273+
await listener.WaitUntilBlockedAsync(cancellation.Token);
274+
275+
var reminderService = initialSilo.ServiceProvider.GetRequiredService<LocalReminderService>();
276+
var barrier = reminderService.TestOnlyWaitForSiloStatusListeners(cancellation.Token);
277+
Assert.False(barrier.IsCompleted);
278+
279+
listener.Release();
280+
joinedSilo = Assert.Single(await startTask.WaitAsync(cancellation.Token));
281+
await barrier;
282+
}
283+
finally
284+
{
285+
listener.Release();
286+
statusOracle.UnSubscribeFromSiloStatusEvents(listener);
287+
if (joinedSilo is null && startTask is not null)
288+
{
289+
using var startupCleanup = new CancellationTokenSource(TestConstants.InitTimeout);
290+
joinedSilo = Assert.Single(await startTask.WaitAsync(startupCleanup.Token));
291+
}
292+
293+
if (joinedSilo is not null)
294+
{
295+
using var siloCleanup = new CancellationTokenSource(TestConstants.InitTimeout);
296+
await fixture.HostedCluster.StopSiloAsync(joinedSilo, siloCleanup.Token);
297+
await fixture.HostedCluster.WaitForLivenessToStabilizeAsync();
298+
}
299+
}
300+
}
301+
257302
[TestSuite("BVT")]
258303
[TestProvider("None")]
259304
[Fact, TestCategory("BVT")]
@@ -546,6 +591,31 @@ public override async ValueTask DisposeAsync()
546591
}
547592
}
548593

594+
private sealed class BlockingSiloStatusListener(SiloAddress localSilo) : ISiloStatusListener
595+
{
596+
private readonly TaskCompletionSource _blocked = new(TaskCreationOptions.RunContinuationsAsynchronously);
597+
private readonly ManualResetEventSlim _release = new();
598+
private int _hasBlocked;
599+
600+
public void SiloStatusChangeNotification(SiloAddress updatedSilo, SiloStatus status)
601+
{
602+
if (updatedSilo.Equals(localSilo)
603+
|| status != SiloStatus.Active
604+
|| Interlocked.Exchange(ref _hasBlocked, 1) != 0)
605+
{
606+
return;
607+
}
608+
609+
_blocked.TrySetResult();
610+
_release.Wait();
611+
}
612+
613+
public Task WaitUntilBlockedAsync(CancellationToken cancellationToken)
614+
=> _blocked.Task.WaitAsync(cancellationToken);
615+
616+
public void Release() => _release.Set();
617+
}
618+
549619
private sealed class NullReturningReminderTable : IReminderTable
550620
{
551621
private int rangeReadCount;

test/Orleans.Testing.Reminders/ReminderServiceLifecycleHarness.cs

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -172,8 +172,12 @@ public async Task WaitForTopologyReconciliationAsync(CancellationToken cancellat
172172
await Task.WhenAll(
173173
_cluster.WaitForLivenessToStabilizeAsync().WaitAsync(cancellationToken),
174174
_cluster.WaitForClusterManifestToStabilizeAsync().WaitAsync(cancellationToken));
175+
var activeSilos = _cluster.GetActiveSilos();
176+
await Task.WhenAll(activeSilos.Select(silo =>
177+
silo.ServiceProvider.GetRequiredService<LocalReminderService>()
178+
.TestOnlyWaitForSiloStatusListeners(cancellationToken)));
175179
await RefreshAsync(cancellationToken);
176-
var barriers = _cluster.GetActiveSilos().Select(silo =>
180+
var barriers = activeSilos.Select(silo =>
177181
silo.ServiceProvider.GetRequiredService<LocalReminderService>()
178182
.TestOnlyWaitForRangeChangeReconciliation(cancellationToken));
179183
await Task.WhenAll(barriers);

0 commit comments

Comments
 (0)