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
46 changes: 46 additions & 0 deletions src/Orleans.Runtime/Diagnostics/GrainDirectoryEvents.cs
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,52 @@ internal sealed class RangeOperationCompleted(
public readonly bool Canceled = canceled;
}

internal sealed class MembershipVersionApplied(
SiloAddress siloAddress,
MembershipVersion version) : GrainDirectoryEvent(siloAddress, partitionIndex: -1, version, RingRange.Empty);

internal sealed class MembershipVersionObserved(
SiloAddress siloAddress,
int partitionIndex,
MembershipVersion version) : GrainDirectoryEvent(siloAddress, partitionIndex, version, RingRange.Empty);

internal static void EmitMembershipVersionApplied(
SiloAddress siloAddress,
MembershipVersion version)
{
if (!Listener.IsEnabled(nameof(MembershipVersionApplied)))
{
return;
}

Emit(siloAddress, version);

[MethodImpl(MethodImplOptions.NoInlining)]
static void Emit(SiloAddress siloAddress, MembershipVersion version)
{
Listener.Write(nameof(MembershipVersionApplied), new MembershipVersionApplied(siloAddress, version));
}
}

internal static void EmitMembershipVersionObserved(
SiloAddress siloAddress,
int partitionIndex,
MembershipVersion version)
{
if (!Listener.IsEnabled(nameof(MembershipVersionObserved)))
{
return;
}

Emit(siloAddress, partitionIndex, version);

[MethodImpl(MethodImplOptions.NoInlining)]
static void Emit(SiloAddress siloAddress, int partitionIndex, MembershipVersion version)
{
Listener.Write(nameof(MembershipVersionObserved), new MembershipVersionObserved(siloAddress, partitionIndex, version));
}
}

internal static void EmitRangeOperationStarted(
SiloAddress siloAddress,
int partitionIndex,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -351,6 +351,7 @@ private void ProcessMembershipUpdate(DirectoryMembershipSnapshot current)
}

_viewUpdates.Publish(current);
GrainDirectoryEvents.EmitMembershipVersionObserved(_id, _partitionIndex, current.Version);
}

private async Task ReleaseRangeAsync(DirectoryMembershipSnapshot previous, DirectoryMembershipSnapshot current, RingRange removedRange)
Expand Down
2 changes: 2 additions & 0 deletions src/Orleans.Runtime/GrainDirectory/LocalGrainDirectory.cs
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@
using Orleans.GrainDirectory;
using Orleans.Internal;
using Orleans.Runtime.Internal;
using Orleans.Runtime.Diagnostics;
using Orleans.Runtime.Scheduler;

namespace Orleans.Runtime.GrainDirectory
Expand Down Expand Up @@ -259,6 +260,7 @@ void ApplyMembershipSnapshotCore()

appliedClusterMembershipSnapshot = snapshot;
hasAppliedClusterMembershipSnapshot = true;
GrainDirectoryEvents.EmitMembershipVersionApplied(MyAddress, snapshot.Version);
}
}

Expand Down
211 changes: 211 additions & 0 deletions src/Orleans.TestingHost/GrainDirectoryObserver.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,211 @@
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Options;
using Orleans.Configuration;
using Orleans.Runtime;
using Orleans.Runtime.Diagnostics;
using Orleans.Runtime.GrainDirectory;

namespace Orleans.TestingHost;

internal sealed class GrainDirectoryObserver : IObserver<GrainDirectoryEvents.GrainDirectoryEvent>, IDisposable
{
private readonly object _lock = new();
private readonly Dictionary<SiloAddress, MembershipVersion> _localVersions = [];
private readonly Dictionary<(SiloAddress SiloAddress, int PartitionIndex), MembershipVersion> _distributedVersions = [];
private readonly HashSet<RangeOperation> _pendingRangeOperations = [];
private readonly IDisposable _subscription;
private TaskCompletionSource _changed = CreateCompletion();
private Exception? _error;

public GrainDirectoryObserver()
{
_subscription = GrainDirectoryEvents.AllEvents.Subscribe(this);
}

public static bool CanObserve(IReadOnlyCollection<InProcessSiloHandle> activeSilos) =>
activeSilos.All(static silo => TryCreateTarget(silo, out _));

public async Task<bool> WaitForConvergenceAsync(
IReadOnlyCollection<InProcessSiloHandle> activeSilos,
TimeSpan timeout)
{
var targets = activeSilos.Select(CreateTarget).ToArray();
using var timeoutSource = new CancellationTokenSource(timeout);
while (true)
{
Task changed;
lock (_lock)
{
if (_error is { } error)
{
throw new InvalidOperationException("An error occurred while observing grain directory events.", error);
}

if (HasConverged(targets))
{
return true;
}

changed = _changed.Task;
}

try
{
await changed.WaitAsync(timeoutSource.Token);
}
catch (OperationCanceledException) when (timeoutSource.IsCancellationRequested)
{
return false;
}
}
}

public void OnNext(GrainDirectoryEvents.GrainDirectoryEvent value)
{
TaskCompletionSource changed;
lock (_lock)
{
switch (value)
{
case GrainDirectoryEvents.MembershipVersionApplied applied:
UpdateVersion(_localVersions, applied.SiloAddress, applied.Version);
break;
case GrainDirectoryEvents.MembershipVersionObserved observed:
UpdateVersion(
_distributedVersions,
(observed.SiloAddress, observed.PartitionIndex),
observed.Version);
break;
case GrainDirectoryEvents.RangeOperationStarted started:
_pendingRangeOperations.Add(RangeOperation.From(started));
break;
case GrainDirectoryEvents.RangeOperationCompleted completed:
_pendingRangeOperations.Remove(RangeOperation.From(completed));
break;
}

changed = _changed;
_changed = CreateCompletion();
}

changed.TrySetResult();
}

public void OnError(Exception error)
{
TaskCompletionSource changed;
lock (_lock)
{
_error = error;
changed = _changed;
_changed = CreateCompletion();
}

changed.TrySetResult();
}

public void OnCompleted()
{
}

public void Dispose() => _subscription.Dispose();

private bool HasConverged(Target[] targets)
{
foreach (var target in targets)
{
if (!_localVersions.TryGetValue(target.SiloAddress, out var localVersion)
|| localVersion < target.Version)
{
return false;
}

if (target.DistributedPartitionCount > 0)
{
for (var partitionIndex = 0; partitionIndex < target.DistributedPartitionCount; partitionIndex++)
{
if (!_distributedVersions.TryGetValue((target.SiloAddress, partitionIndex), out var version)
|| version < target.Version)
{
return false;
}
}

if (_pendingRangeOperations.Any(operation => operation.SiloAddress.Equals(target.SiloAddress)))
{
return false;
}
}
}

return true;
}

private static Target CreateTarget(InProcessSiloHandle silo)
{
if (TryCreateTarget(silo, out var target))
{
return target;
}

throw new InvalidOperationException(
$"The default grain directory on silo {silo.SiloAddress} does not emit grain directory convergence events.");
}

private static bool TryCreateTarget(InProcessSiloHandle silo, out Target target)
{
var services = silo.ServiceProvider;
var membershipVersion = services.GetRequiredService<IClusterMembershipService>().CurrentSnapshot.Version;
var defaultDirectory = services.GetRequiredService<GrainDirectoryResolver>().DefaultGrainDirectory;
var distributedPartitionCount = defaultDirectory switch
{
null => 0,
DistributedGrainDirectory => services.GetRequiredService<IOptions<GrainDirectoryOptions>>().Value.PartitionsPerSilo,
_ => -1
};

target = new(silo.SiloAddress, membershipVersion, distributedPartitionCount);
return distributedPartitionCount >= 0;
}

private static void UpdateVersion<TKey>(
Dictionary<TKey, MembershipVersion> versions,
TKey key,
MembershipVersion version)
where TKey : notnull
{
if (!versions.TryGetValue(key, out var current) || version > current)
{
versions[key] = version;
}
}

private static TaskCompletionSource CreateCompletion() =>
new(TaskCreationOptions.RunContinuationsAsynchronously);

private readonly record struct Target(
SiloAddress SiloAddress,
MembershipVersion Version,
int DistributedPartitionCount);

private readonly record struct RangeOperation(
SiloAddress SiloAddress,
int PartitionIndex,
MembershipVersion Version,
RingRange Range,
string OperationName)
{
public static RangeOperation From(GrainDirectoryEvents.RangeOperationEvent operation) =>
new(
operation.SiloAddress,
operation.PartitionIndex,
operation.Version,
operation.Range,
operation.OperationName);
}
}
14 changes: 13 additions & 1 deletion src/Orleans.TestingHost/InProcTestCluster.cs
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ public sealed class InProcessTestCluster : IDisposable, IAsyncDisposable
private readonly List<InProcessSiloHandle> _silos = [];
private readonly StringBuilder _log = new();
private readonly InMemoryTransportConnectionHub _transportHub = new();
private readonly GrainDirectoryObserver _grainDirectoryObserver = new();
private readonly InProcessGrainDirectory _grainDirectory;
private readonly InProcessMembershipTable _membershipTable;
private bool _disposed;
Expand Down Expand Up @@ -365,8 +366,17 @@ public async Task WaitForLivenessToStabilizeAsync(bool didKill = false)
var activeSilos = GetActiveSilos().ToArray();
var testHooks = activeSilos.Select(static silo => (ITestHooks)silo.ServiceProvider.GetRequiredService<TestHooksSystemTarget>()).ToArray();
var gatewayManager = Client.ServiceProvider.GetRequiredService<GatewayManager>();
Func<TimeSpan, Task<bool>>? waitForGrainDirectoryConvergence =
GrainDirectoryObserver.CanObserve(activeSilos)
? timeout => _grainDirectoryObserver.WaitForConvergenceAsync(activeSilos, timeout)
: null;
WriteLog(Environment.NewLine + Environment.NewLine + "WaitForLivenessToStabilize is waiting up to {0} for {1} active silo(s)", stabilizationTime, activeSilos.Length);
if (await LivenessStabilizationHelper.WaitForExpectedActiveSilosAndGatewaysAsync(activeSilos, testHooks, gatewayManager, stabilizationTime))
if (await LivenessStabilizationHelper.WaitForExpectedActiveSilosAndGatewaysAsync(
activeSilos,
testHooks,
gatewayManager,
stabilizationTime,
waitForGrainDirectoryConvergence))
{
WriteLog("WaitForLivenessToStabilize observed stable active silo and gateway views");
}
Expand Down Expand Up @@ -918,6 +928,7 @@ await Task.Run(async () =>
ClientHost = null;

PortAllocator?.Dispose();
_grainDirectoryObserver.Dispose();
});

_disposed = true;
Expand All @@ -938,6 +949,7 @@ public void Dispose()

ClientHost?.Dispose();
PortAllocator?.Dispose();
_grainDirectoryObserver.Dispose();

_disposed = true;
}
Expand Down
27 changes: 17 additions & 10 deletions src/Orleans.TestingHost/LivenessStabilizationHelper.cs
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,8 @@ public static async Task<bool> WaitForExpectedActiveSilosAndGatewaysAsync(
IReadOnlyCollection<SiloHandle> activeSilos,
IReadOnlyCollection<ITestHooks> testHooks,
GatewayManager gatewayManager,
TimeSpan timeout)
TimeSpan timeout,
Func<TimeSpan, Task<bool>>? waitForGrainDirectoryConvergence = null)
{
ArgumentNullException.ThrowIfNull(gatewayManager);

Expand All @@ -27,10 +28,15 @@ public static async Task<bool> WaitForExpectedActiveSilosAndGatewaysAsync(
return false;
}

var remaining = timeout - stopwatch.Elapsed;
return remaining <= TimeSpan.Zero
? ActiveGatewaysMatch(gatewayManager, activeSilos)
: await WaitForExpectedActiveGatewaysAsync(activeSilos, gatewayManager, remaining);
var remaining = GetRemainingTime(timeout, stopwatch.Elapsed);
if (waitForGrainDirectoryConvergence is not null
&& !await waitForGrainDirectoryConvergence(remaining))
{
return false;
}

remaining = GetRemainingTime(timeout, stopwatch.Elapsed);
return await WaitForExpectedActiveGatewaysAsync(activeSilos, gatewayManager, remaining);
}

public static async Task<bool> WaitForExpectedActiveSilosAsync(
Expand Down Expand Up @@ -60,6 +66,12 @@ public static async Task<bool> WaitForExpectedActiveSilosAsync(
}
}

private static TimeSpan GetRemainingTime(TimeSpan timeout, TimeSpan elapsed)
{
var remaining = timeout - elapsed;
return remaining > TimeSpan.Zero ? remaining : TimeSpan.Zero;
}

private static async Task<bool> WaitForExpectedActiveGatewaysAsync(
IReadOnlyCollection<SiloHandle> activeSilos,
GatewayManager gatewayManager,
Expand Down Expand Up @@ -97,11 +109,6 @@ private static async Task<bool> WaitForExpectedActiveGatewaysAsync(
}
}

private static bool ActiveGatewaysMatch(GatewayManager gatewayManager, IReadOnlyCollection<SiloHandle> activeSilos)
{
return GatewaysMatch(gatewayManager.GetLiveGateways(), GetExpectedGatewayAddresses(activeSilos));
}

private static HashSet<SiloAddress> GetExpectedGatewayAddresses(IReadOnlyCollection<SiloHandle> activeSilos)
{
return activeSilos
Expand Down
Loading
Loading