diff --git a/src/Orleans.Clustering.ZooKeeper/Orleans.Clustering.ZooKeeper.csproj b/src/Orleans.Clustering.ZooKeeper/Orleans.Clustering.ZooKeeper.csproj index 806eefb0e86..4d5467dc510 100644 --- a/src/Orleans.Clustering.ZooKeeper/Orleans.Clustering.ZooKeeper.csproj +++ b/src/Orleans.Clustering.ZooKeeper/Orleans.Clustering.ZooKeeper.csproj @@ -11,6 +11,7 @@ + diff --git a/src/Orleans.Clustering.ZooKeeper/ZooKeeperBasedMembershipTable.cs b/src/Orleans.Clustering.ZooKeeper/ZooKeeperBasedMembershipTable.cs index f81f18834fd..18c83e15b6a 100644 --- a/src/Orleans.Clustering.ZooKeeper/ZooKeeperBasedMembershipTable.cs +++ b/src/Orleans.Clustering.ZooKeeper/ZooKeeperBasedMembershipTable.cs @@ -12,6 +12,7 @@ using Microsoft.Extensions.Options; using Orleans.Configuration; using Orleans.Runtime.Host; +using Polly; namespace Orleans.Runtime.Membership { @@ -40,9 +41,14 @@ public partial class ZooKeeperBasedMembershipTable : IMembershipTable { private readonly ILogger logger; - private const int ZOOKEEPER_SESSION_TIMEOUT = 10_000; + internal const int ZOOKEEPER_SESSION_TIMEOUT = 10_000; + internal const int MAX_MEMBERSHIP_SNAPSHOT_ATTEMPTS = 5; + internal const int MAX_CLEANUP_ROW_ATTEMPTS = 5; private readonly ZooKeeperWatcher watcher; + private readonly Func _createSession; + private readonly ResiliencePipeline _readRetryPipeline; + private readonly Action? _observeOperation; /// /// The deployment connection string. for eg. "192.168.1.1,192.168.1.2/ClusterId" @@ -73,6 +79,17 @@ public ZooKeeperBasedMembershipTable( ILogger logger, IOptions membershipTableOptions, IOptions clusterOptions) + : this(logger, membershipTableOptions, clusterOptions, null, null) + { + } + + internal ZooKeeperBasedMembershipTable( + ILogger logger, + IOptions membershipTableOptions, + IOptions clusterOptions, + Func? createSession, + ResiliencePipeline? readRetryPipeline, + Action? observeOperation = null) { ArgumentNullException.ThrowIfNull(logger); ArgumentNullException.ThrowIfNull(membershipTableOptions); @@ -84,6 +101,9 @@ public ZooKeeperBasedMembershipTable( this.clusterPath = "/" + clusterOptions.Value.ClusterId; rootConnectionString = options.ConnectionString; deploymentConnectionString = options.ConnectionString + this.clusterPath; + _createSession = createSession ?? (readOnly => CreateSession(deploymentConnectionString, watcher, readOnly)); + _readRetryPipeline = readRetryPipeline ?? ZooKeeperReadRetryPolicy.CreatePipeline(logger, TimeProvider.System); + _observeOperation = observeOperation; } /// @@ -128,10 +148,16 @@ await UsingZookeeper(rootConnectionString, async zk => public Task ReadRow(SiloAddress siloAddress) => ReadRowAsync(siloAddress, CancellationToken.None); /// + /// + /// Connection-loss failures in native reads are retried up to four times on the operation's session. + /// Before each retry, the operation waits up to the session timeout for that session's next connected event. + /// When the wait expires or the session becomes terminal, the retry proceeds and preserves the native outcome. + /// The table and child versions fence each complete snapshot pass. Concurrent canonical + /// modifications restart the pass up to five total attempts. + /// public Task ReadRowAsync(SiloAddress siloAddress, CancellationToken cancellationToken = default) { - return UsingZookeeper(zk => ReadCoreAsync(zk, siloAddress, cancellationToken), - this.deploymentConnectionString, this.watcher, cancellationToken, canBeReadOnly: true); + return ReadAsync(() => _createSession(true), _readRetryPipeline, siloAddress, cancellationToken); } /// @@ -145,15 +171,48 @@ public Task ReadRowAsync(SiloAddress siloAddress, Cancellat public Task ReadAll() => ReadAllAsync(CancellationToken.None); /// + /// + /// Rows are read sequentially on an operation-owned connection. + /// Each membership record is read before its heartbeat. + /// Table and child-version checks fence the complete snapshot. Concurrent canonical + /// modifications restart the complete sequential pass up to five total attempts. + /// Connection-loss failures in native reads are retried up to four times on the same session. + /// Before each retry, the operation waits up to the session timeout for that session's next connected event. + /// When the wait expires or the session becomes terminal, the retry proceeds and preserves the native outcome. + /// Caller cancellation stops further requests while admitted requests and client close complete. + /// public Task ReadAllAsync(CancellationToken cancellationToken = default) { - return ReadAllAsync(this.deploymentConnectionString, this.watcher, cancellationToken); + return ReadAsync(() => _createSession(true), _readRetryPipeline, null, cancellationToken); } - internal static Task ReadAllAsync(string deploymentConnectionString, ZooKeeperWatcher watcher, CancellationToken cancellationToken) + internal static Task ReadAsync( + Func createSession, + ResiliencePipeline pipeline, + SiloAddress? siloAddress, + CancellationToken cancellationToken) { - return UsingZookeeper(zk => ReadCoreAsync(zk, null, cancellationToken), - deploymentConnectionString, watcher, cancellationToken, canBeReadOnly: true); + return ZooKeeperSession.ExecuteSessionAsync(createSession, + session => ReadCoreAsync( + ZooKeeperReadRetryPolicy.Wrap( + session.Operations, + pipeline, + cancellationToken, + session.ConnectionMonitor), + siloAddress, + cancellationToken), + cancellationToken); + } + + internal static ZooKeeperSession CreateSession(string connectionString, ZooKeeperWatcher watcher, bool readOnly) + { + var sessionWatcher = watcher.CreateSessionWatcher(); + var client = new ZooKeeper(connectionString, ZOOKEEPER_SESSION_TIMEOUT, sessionWatcher, readOnly); + return new ZooKeeperSession( + new NativeOperations(path => client.getDataAsync(path), path => client.getChildrenAsync(path), + client.sync, operations => client.multiAsync(operations), client.setDataAsync), + client.closeAsync, + sessionWatcher); } internal static async Task ReadCoreAsync( @@ -162,7 +221,7 @@ internal static async Task ReadCoreAsync( cancellationToken.ThrowIfCancellationRequested(); await zk.Sync("/"); // Retries retain this session's ordered view. - while (true) + for (var attempt = 0; attempt < MAX_MEMBERSHIP_SNAPSHOT_ATTEMPTS; attempt++) { cancellationToken.ThrowIfCancellationRequested(); Stat before; @@ -179,38 +238,40 @@ internal static async Task ReadCoreAsync( addresses = [siloAddress]; } - var pendingRows = Task.WhenAll(addresses.Select(address => GetRow(zk, address, siloAddress is not null, cancellationToken))); - Tuple?[] rows; - try + var rows = new List>(); + KeeperException.NoNodeException? missingRow = null; + foreach (var address in addresses) { - rows = await pendingRows; - } - catch (KeeperException.NoNodeException) - { - // Observe every parallel read: a removed row must not hide another request's failure. - var failure = pendingRows.Exception!.InnerExceptions.FirstOrDefault(exception => exception is not KeeperException.NoNodeException); - if (failure is not null) + cancellationToken.ThrowIfCancellationRequested(); + try { - System.Runtime.ExceptionServices.ExceptionDispatchInfo.Capture(failure).Throw(); + if (await GetRow(zk, address, siloAddress is not null, cancellationToken) is { } row) + { + rows.Add(row); + } } - - cancellationToken.ThrowIfCancellationRequested(); - var current = await zk.GetData("/"); - if (SameVersion(before, current.Stat)) + catch (KeeperException.NoNodeException exception) { - throw; + missingRow = exception; + break; } - - continue; } cancellationToken.ThrowIfCancellationRequested(); var after = await zk.GetData("/"); if (SameVersion(before, after.Stat)) { - return new MembershipTableData(rows.OfType>().ToList(), ConvertToTableVersion(after.Stat)); + if (missingRow is not null) + { + System.Runtime.ExceptionServices.ExceptionDispatchInfo.Capture(missingRow).Throw(); + } + + return new MembershipTableData(rows, ConvertToTableVersion(after.Stat)); } } + + throw new OrleansException( + $"Unable to read a consistent ZooKeeper membership snapshot after {MAX_MEMBERSHIP_SNAPSHOT_ATTEMPTS} attempts."); } private static bool SameVersion(Stat before, Stat after) => @@ -238,14 +299,18 @@ private static bool SameVersion(Stat before, Stat after) => public Task InsertRow(MembershipEntry entry, TableVersion tableVersion) => InsertRowAsync(entry, tableVersion, CancellationToken.None); /// + /// + /// The conditional transaction is submitted once. Connection loss propagates when its + /// commit outcome is unknown, preserving the caller's ability to resolve that outcome. + /// public Task InsertRowAsync(MembershipEntry entry, TableVersion tableVersion, CancellationToken cancellationToken = default) { ArgumentNullException.ThrowIfNull(entry); ArgumentNullException.ThrowIfNull(tableVersion); cancellationToken.ThrowIfCancellationRequested(); - return UsingZookeeper(zk => InsertRowCoreAsync(zk, entry, tableVersion, cancellationToken), - this.deploymentConnectionString, this.watcher, cancellationToken); + return ZooKeeperSession.ExecuteAsync(() => _createSession(false), + zk => InsertRowCoreAsync(zk, entry, tableVersion, cancellationToken), cancellationToken); } internal static async Task InsertRowCoreAsync( @@ -300,6 +365,10 @@ await zk.Multi( public Task UpdateRow(MembershipEntry entry, string etag, TableVersion tableVersion) => UpdateRowAsync(entry, etag, tableVersion, CancellationToken.None); /// + /// + /// The conditional transaction is submitted once. Connection loss propagates when its + /// commit outcome is unknown, preserving the caller's ability to resolve that outcome. + /// public Task UpdateRowAsync(MembershipEntry entry, string etag, TableVersion tableVersion, CancellationToken cancellationToken = default) { ArgumentNullException.ThrowIfNull(entry); @@ -307,8 +376,8 @@ public Task UpdateRowAsync(MembershipEntry entry, string etag, TableVersio ArgumentNullException.ThrowIfNull(tableVersion); cancellationToken.ThrowIfCancellationRequested(); - return UsingZookeeper(zk => UpdateRowCoreAsync(zk, entry, etag, tableVersion, cancellationToken), - this.deploymentConnectionString, this.watcher, cancellationToken); + return ZooKeeperSession.ExecuteAsync(() => _createSession(false), + zk => UpdateRowCoreAsync(zk, entry, etag, tableVersion, cancellationToken), cancellationToken); } internal static async Task UpdateRowCoreAsync( @@ -431,7 +500,7 @@ internal sealed class NativeOperations( internal Func> SetData { get; } = setData; } - private static async Task UsingZookeeper(Func> zkMethod, string deploymentConnectionString, ZooKeeperWatcher watcher, CancellationToken cancellationToken, bool canBeReadOnly = false) + private async Task UsingZookeeper(Func> zkMethod, string deploymentConnectionString, ZooKeeperWatcher watcher, CancellationToken cancellationToken, bool canBeReadOnly = false) { cancellationToken.ThrowIfCancellationRequested(); var operation = ZooKeeper.Using(deploymentConnectionString, ZOOKEEPER_SESSION_TIMEOUT, watcher, zk => @@ -442,13 +511,15 @@ private static async Task UsingZookeeper(Func> z operations => zk.multiAsync(operations), zk.setDataAsync)); }, canBeReadOnly); - return await AwaitOperationAsync(operation, cancellationToken); + return await AwaitOperationAsync(operation, cancellationToken, _observeOperation); } - internal static async Task AwaitOperationAsync(Task operation, CancellationToken cancellationToken) + internal static async Task AwaitOperationAsync( + Task operation, CancellationToken cancellationToken, Action? observeOperation = null) { // ZooKeeperNetEx is tokenless. Keep the client alive until pending requests and // asynchronous disposal finish, observing failures even if the caller stops waiting. + observeOperation?.Invoke(operation); operation.Ignore(); return await operation.WaitAsync(cancellationToken); } @@ -462,6 +533,7 @@ private async Task UsingZookeeper(string connectString, Func zk return zkMethod(zk); }); + _observeOperation?.Invoke(operation); operation.Ignore(); await operation.WaitAsync(cancellationToken); } @@ -520,6 +592,11 @@ internal static T Deserialize(byte[] data) public Task CleanupDefunctSiloEntries(DateTimeOffset beforeDate) => CleanupDefunctSiloEntriesAsync(beforeDate, CancellationToken.None); /// + /// + /// Rows are evaluated sequentially in the order returned by ZooKeeper. A row whose + /// conditional delete conflicts is re-evaluated up to five total attempts before + /// the operation reports contention. + /// public Task CleanupDefunctSiloEntriesAsync(DateTimeOffset beforeDate, CancellationToken cancellationToken = default) { return UsingZookeeper(zk => CleanupCoreAsync(zk, beforeDate, cancellationToken), @@ -531,14 +608,19 @@ internal static async Task CleanupCoreAsync(NativeOperations zk, DateTimeO cancellationToken.ThrowIfCancellationRequested(); var children = await zk.GetChildren("/"); var cutoff = beforeDate.UtcDateTime; - await Task.WhenAll(children.Children.Select(child => CleanupRowAsync(zk, "/" + child, cutoff, cancellationToken))); + foreach (var child in children.Children) + { + cancellationToken.ThrowIfCancellationRequested(); + await CleanupRowAsync(zk, "/" + child, cutoff, cancellationToken); + } + return true; } private static async Task CleanupRowAsync(NativeOperations zk, string rowPath, DateTime cutoff, CancellationToken cancellationToken) { var heartbeatPath = rowPath + "/IAmAlive"; - while (true) + for (var attempt = 0; attempt < MAX_CLEANUP_ROW_ATTEMPTS; attempt++) { cancellationToken.ThrowIfCancellationRequested(); DataResult row; @@ -599,6 +681,9 @@ await zk.Multi( throw; } } + + throw new OrleansException( + $"Unable to clean ZooKeeper membership row '{rowPath}' after {MAX_CLEANUP_ROW_ATTEMPTS} concurrent modifications."); } [LoggerMessage( @@ -615,23 +700,143 @@ await zk.Multi( } /// - /// the state of every ZooKeeper client and its push notifications are published using watchers. - /// in orleans the watcher is only for debugging purposes + /// Publishes ZooKeeper connection transitions for same-session read retry coordination + /// and logs watcher events for diagnostics. /// - internal partial class ZooKeeperWatcher : Watcher + internal partial class ZooKeeperWatcher : Watcher, IZooKeeperConnectionMonitor { private readonly ILogger logger; + private readonly TimeProvider _timeProvider; + private readonly TimeSpan _reconnectTimeout; + private readonly object _connectionLock = new(); + private TaskCompletionSource _connectionChanged = NewConnectionChangedSource(); + private long _connectedGeneration; + private bool _connected; + private bool _terminal; + public ZooKeeperWatcher(ILogger logger) + : this(logger, TimeProvider.System, TimeSpan.FromMilliseconds(ZooKeeperBasedMembershipTable.ZOOKEEPER_SESSION_TIMEOUT)) + { + } + + internal ZooKeeperWatcher(ILogger logger, TimeProvider timeProvider, TimeSpan reconnectTimeout) { this.logger = logger; + _timeProvider = timeProvider; + _reconnectTimeout = reconnectTimeout; } + internal ZooKeeperWatcher CreateSessionWatcher() => new(logger, _timeProvider, _reconnectTimeout); + public override Task process(WatchedEvent @event) { + ProcessConnectionState(@event.getState()); LogDebugWatchedEvent(@event); return Task.CompletedTask; } + internal void ProcessConnectionState(Event.KeeperState state) + { + TaskCompletionSource? changed = null; + lock (_connectionLock) + { + switch (state) + { + case Event.KeeperState.SyncConnected: + case Event.KeeperState.ConnectedReadOnly: + _connected = true; + _connectedGeneration++; + changed = _connectionChanged; + _connectionChanged = NewConnectionChangedSource(); + break; + case Event.KeeperState.AuthFailed: + case Event.KeeperState.Expired: + _connected = false; + _terminal = true; + changed = _connectionChanged; + _connectionChanged = NewConnectionChangedSource(); + break; + case Event.KeeperState.Disconnected: + _connected = false; + changed = _connectionChanged; + _connectionChanged = NewConnectionChangedSource(); + break; + } + } + + changed?.TrySetResult(); + } + + long IZooKeeperConnectionMonitor.CaptureAttemptGeneration() + { + lock (_connectionLock) + { + // Requests admitted while connecting execute on the next connection. + return _connected ? _connectedGeneration : _connectedGeneration + 1; + } + } + + void IZooKeeperConnectionMonitor.ReportConnectionLoss(long connectedGeneration) + { + TaskCompletionSource? changed = null; + lock (_connectionLock) + { + if (_connected && _connectedGeneration <= connectedGeneration) + { + _connected = false; + changed = _connectionChanged; + _connectionChanged = NewConnectionChangedSource(); + } + } + + changed?.TrySetResult(); + } + + async ValueTask IZooKeeperConnectionMonitor.WaitForConnectionAfterAsync( + long connectedGeneration, + CancellationToken cancellationToken) + { + using var timeoutCancellation = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + var timeout = Task.Delay(_reconnectTimeout, _timeProvider, timeoutCancellation.Token); + try + { + while (true) + { + Task changed; + lock (_connectionLock) + { + if (_connected && _connectedGeneration > connectedGeneration) + { + return true; + } + + if (_terminal) + { + return false; + } + + changed = _connectionChanged.Task; + } + + var completed = await Task.WhenAny(changed, timeout); + cancellationToken.ThrowIfCancellationRequested(); + if (completed == timeout) + { + return false; + } + + await changed; + } + } + finally + { + timeoutCancellation.Cancel(); + } + } + + private static TaskCompletionSource NewConnectionChangedSource() => + new(TaskCreationOptions.RunContinuationsAsynchronously); + [LoggerMessage( Level = LogLevel.Debug, Message = "{EventString}" diff --git a/src/Orleans.Clustering.ZooKeeper/ZooKeeperConnectionMonitor.cs b/src/Orleans.Clustering.ZooKeeper/ZooKeeperConnectionMonitor.cs new file mode 100644 index 00000000000..a1e91dda530 --- /dev/null +++ b/src/Orleans.Clustering.ZooKeeper/ZooKeeperConnectionMonitor.cs @@ -0,0 +1,14 @@ +using System; +using System.Threading; +using System.Threading.Tasks; + +namespace Orleans.Runtime.Membership; + +internal interface IZooKeeperConnectionMonitor +{ + long CaptureAttemptGeneration(); + + ValueTask WaitForConnectionAfterAsync(long connectedGeneration, CancellationToken cancellationToken); + + void ReportConnectionLoss(long connectedGeneration); +} diff --git a/src/Orleans.Clustering.ZooKeeper/ZooKeeperGatewayListProvider.cs b/src/Orleans.Clustering.ZooKeeper/ZooKeeperGatewayListProvider.cs index aa9cebecda2..c35bd695c68 100644 --- a/src/Orleans.Clustering.ZooKeeper/ZooKeeperGatewayListProvider.cs +++ b/src/Orleans.Clustering.ZooKeeper/ZooKeeperGatewayListProvider.cs @@ -7,6 +7,7 @@ using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using Orleans.Configuration; +using Polly; namespace Orleans.Runtime.Membership { @@ -27,6 +28,8 @@ public class ZooKeeperGatewayListProvider : IGatewayListProvider /// private readonly string _deploymentConnectionString; private readonly TimeSpan _maxStaleness; + private readonly Func _createSession; + private readonly ResiliencePipeline _readRetryPipeline; /// /// Initializes a new instance of the class. @@ -44,6 +47,17 @@ public ZooKeeperGatewayListProvider( IOptions options, IOptions gatewayOptions, IOptions clusterOptions) + : this(logger, options, gatewayOptions, clusterOptions, null, null) + { + } + + internal ZooKeeperGatewayListProvider( + ILogger logger, + IOptions options, + IOptions gatewayOptions, + IOptions clusterOptions, + Func? createSession, + ResiliencePipeline? readRetryPipeline) { ArgumentNullException.ThrowIfNull(logger); ArgumentNullException.ThrowIfNull(options); @@ -54,6 +68,8 @@ public ZooKeeperGatewayListProvider( _deploymentPath = "/" + clusterOptions.Value.ClusterId; _deploymentConnectionString = options.Value.ConnectionString + _deploymentPath; _maxStaleness = gatewayOptions.Value.GatewayListRefreshPeriod; + _createSession = createSession ?? (() => ZooKeeperBasedMembershipTable.CreateSession(_deploymentConnectionString, _watcher, true)); + _readRetryPipeline = readRetryPipeline ?? ZooKeeperReadRetryPolicy.CreatePipeline(logger, TimeProvider.System); } /// @@ -65,9 +81,13 @@ public ZooKeeperGatewayListProvider( /// Returns the list of gateways (silos) that can be used by a client to connect to Orleans cluster. /// The Uri is in the form of: "gwy.tcp://IP:port/Generation". See Utils.ToGatewayUri and Utils.ToSiloAddress for more details about Uri format. /// + /// + /// Gateway discovery uses a version-fenced membership snapshot. Native reads retry connection-loss + /// failures up to four times on the operation's session before propagating the final failure. + /// public async Task> GetGateways() { - var membershipTableData = await ZooKeeperBasedMembershipTable.ReadAllAsync(this._deploymentConnectionString, this._watcher, CancellationToken.None); + var membershipTableData = await ZooKeeperBasedMembershipTable.ReadAsync(_createSession, _readRetryPipeline, null, CancellationToken.None); return membershipTableData.Members.Select(e => e.Item1). Where(m => m.Status == SiloStatus.Active && m.ProxyPort != 0). Select(m => diff --git a/src/Orleans.Clustering.ZooKeeper/ZooKeeperReadRetryPolicy.cs b/src/Orleans.Clustering.ZooKeeper/ZooKeeperReadRetryPolicy.cs new file mode 100644 index 00000000000..a9ac4f6243b --- /dev/null +++ b/src/Orleans.Clustering.ZooKeeper/ZooKeeperReadRetryPolicy.cs @@ -0,0 +1,96 @@ +using System; +using System.Threading; +using System.Threading.Tasks; +using Microsoft.Extensions.Logging; +using org.apache.zookeeper; +using Polly; +using Polly.Retry; + +namespace Orleans.Runtime.Membership; + +internal static partial class ZooKeeperReadRetryPolicy +{ + internal const int MaxRetryAttempts = 4; + internal static readonly TimeSpan RetryDelay = TimeSpan.FromMilliseconds(250); + + internal static ResiliencePipeline CreatePipeline(ILogger logger, TimeProvider timeProvider) => + new ResiliencePipelineBuilder { TimeProvider = timeProvider } + .AddRetry(new RetryStrategyOptions + { + MaxRetryAttempts = MaxRetryAttempts, + BackoffType = DelayBackoffType.Exponential, + Delay = RetryDelay, + ShouldHandle = new PredicateBuilder().Handle(), + OnRetry = args => + { + LogWarningRetryRead(logger, args.Outcome.Exception!, args.Context.OperationKey!, + args.AttemptNumber + 1, MaxRetryAttempts, args.RetryDelay.TotalMilliseconds); + return default; + } + }) + .Build(); + + internal static ZooKeeperBasedMembershipTable.NativeOperations Wrap( + ZooKeeperBasedMembershipTable.NativeOperations native, + ResiliencePipeline pipeline, + CancellationToken cancellationToken, + IZooKeeperConnectionMonitor? connectionMonitor = null) => + new( + path => ExecuteAsync("GetData", () => native.GetData(path), pipeline, cancellationToken, connectionMonitor), + path => ExecuteAsync("GetChildren", () => native.GetChildren(path), pipeline, cancellationToken, connectionMonitor), + path => ExecuteAsync("Sync", async () => + { + await native.Sync(path); + return true; + }, pipeline, cancellationToken, connectionMonitor), + native.Multi, + native.SetData); + + private static async Task ExecuteAsync( + string operationName, + Func> operation, + ResiliencePipeline pipeline, + CancellationToken cancellationToken, + IZooKeeperConnectionMonitor? connectionMonitor) + { + var context = ResilienceContextPool.Shared.Get(operationName, cancellationToken); + long? reconnectAfter = null; + try + { + return await pipeline.ExecuteAsync(async context => + { + context.CancellationToken.ThrowIfCancellationRequested(); + if (reconnectAfter is { } connectedGeneration && connectionMonitor is not null) + { + await connectionMonitor.WaitForConnectionAfterAsync( + connectedGeneration, + context.CancellationToken); + } + + context.CancellationToken.ThrowIfCancellationRequested(); + var attemptGeneration = connectionMonitor?.CaptureAttemptGeneration() ?? 0; + // Await actual native completion: cancellation ends admission, not an in-flight request. + try + { + return await operation(); + } + catch (KeeperException.ConnectionLossException) + { + connectionMonitor?.ReportConnectionLoss(attemptGeneration); + reconnectAfter = attemptGeneration; + throw; + } + }, context); + } + finally + { + ResilienceContextPool.Shared.Return(context); + } + } + + [LoggerMessage( + Level = LogLevel.Warning, + Message = "ZooKeeper {Operation} lost its connection. Retrying in {DelayMilliseconds}ms ({Retry}/{MaxRetries}).")] + private static partial void LogWarningRetryRead( + ILogger logger, Exception exception, string operation, int retry, int maxRetries, double delayMilliseconds); +} diff --git a/src/Orleans.Clustering.ZooKeeper/ZooKeeperSession.cs b/src/Orleans.Clustering.ZooKeeper/ZooKeeperSession.cs new file mode 100644 index 00000000000..d99894dc893 --- /dev/null +++ b/src/Orleans.Clustering.ZooKeeper/ZooKeeperSession.cs @@ -0,0 +1,73 @@ +using System; +using System.Threading; +using System.Threading.Tasks; + +namespace Orleans.Runtime.Membership; + +internal sealed class ZooKeeperSession +{ + private readonly ZooKeeperBasedMembershipTable.NativeOperations _operations; + private readonly Func _close; + // Bind synchronously; the owned task supplies the completion semantics. + private readonly TaskCompletionSource _completion = new(); + + internal ZooKeeperSession( + ZooKeeperBasedMembershipTable.NativeOperations operations, + Func close, + IZooKeeperConnectionMonitor? connectionMonitor = null) + { + _operations = operations; + _close = close; + ConnectionMonitor = connectionMonitor; + Completion = _completion.Task.Unwrap(); + Completion.Ignore(); + } + + internal Task Completion { get; } + internal IZooKeeperConnectionMonitor? ConnectionMonitor { get; } + internal ZooKeeperBasedMembershipTable.NativeOperations Operations => _operations; + + internal static Task ExecuteAsync( + Func createSession, + Func> operation, + CancellationToken cancellationToken) + => ExecuteSessionAsync(createSession, session => operation(session.Operations), cancellationToken); + + internal static Task ExecuteSessionAsync( + Func createSession, + Func> operation, + CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + var session = createSession(); + var completion = session.RunAsync(operation); + session._completion.SetResult(completion); + return ZooKeeperBasedMembershipTable.AwaitOperationAsync(completion, cancellationToken); + } + + private async Task RunAsync(Func> operation) + { + T result; + try + { + result = await operation(this); + } + catch (Exception primary) + { + try + { + await _close(); + } + catch (Exception secondary) + { + throw new AggregateException("The ZooKeeper operation and its close both failed.", primary, secondary); + } + + throw; + } + + // The callback joins native requests before this tokenless close. + await _close(); + return result; + } +} diff --git a/test/Extensions/Orleans.Clustering.ZooKeeper.Tests/Orleans.Clustering.ZooKeeper.Tests.csproj b/test/Extensions/Orleans.Clustering.ZooKeeper.Tests/Orleans.Clustering.ZooKeeper.Tests.csproj index 74ae4ddd3cd..a41b183a5b8 100644 --- a/test/Extensions/Orleans.Clustering.ZooKeeper.Tests/Orleans.Clustering.ZooKeeper.Tests.csproj +++ b/test/Extensions/Orleans.Clustering.ZooKeeper.Tests/Orleans.Clustering.ZooKeeper.Tests.csproj @@ -6,7 +6,9 @@ + + diff --git a/test/Extensions/Orleans.Clustering.ZooKeeper.Tests/ZooKeeperBasedMembershipTableUnitTests.cs b/test/Extensions/Orleans.Clustering.ZooKeeper.Tests/ZooKeeperBasedMembershipTableUnitTests.cs index b33fc088b03..ac0459edced 100644 --- a/test/Extensions/Orleans.Clustering.ZooKeeper.Tests/ZooKeeperBasedMembershipTableUnitTests.cs +++ b/test/Extensions/Orleans.Clustering.ZooKeeper.Tests/ZooKeeperBasedMembershipTableUnitTests.cs @@ -1,7 +1,10 @@ using System; +using System.Collections.Concurrent; using System.Collections.Generic; +using System.Globalization; using System.Linq; using System.Net; +using System.Net.Sockets; using System.Reflection; using System.Threading; using System.Threading.Tasks; @@ -22,6 +25,29 @@ namespace UnitTests.MembershipTests [TestArea("Membership")] public sealed class ZooKeeperBasedMembershipTableUnitTests { + [Fact] + public async Task NativeSocketDiagnostics_ExistingSourceReportsOriginalCompletionError() + { + using var destination = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + destination.Bind(new IPEndPoint(IPAddress.Loopback, 0)); + var messages = new ConcurrentQueue(); + using var diagnostics = new NativeSocketDiagnostics(messages.Enqueue); + Assert.Contains(messages, message => message.EndsWith("Error listener enabled", StringComparison.Ordinal)); + using var client = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + using var operation = new SocketAsyncEventArgs { RemoteEndPoint = destination.LocalEndPoint }; + var completion = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + operation.Completed += (_, result) => completion.TrySetResult(result.SocketError); + + if (!client.ConnectAsync(operation)) + { + completion.TrySetResult(operation.SocketError); + } + + Assert.Equal(SocketError.ConnectionRefused, await completion.Task.WaitAsync(TestContext.Current.CancellationToken)); + Assert.Contains(messages, message => + message.Contains($"Socket#{client.GetHashCode()}; UpdateStatusAfterSocketError; errorCode:ConnectionRefused", StringComparison.Ordinal)); + } + [Fact] public void Constructor_NullLogger_ThrowsArgumentNullException() { @@ -381,20 +407,356 @@ public async Task Read_ConcurrentHeartbeat_PreservesCanonicalFence(bool pointRea } [Fact] - public async Task Read_ParallelMissingRowAndAuthorizationFailure_PropagatesAuthorizationFailure() + public async Task ReadAll_ReadsRowsSequentially_WithMemberBeforeHeartbeat() + { + var (fake, first) = await CreateNativeTable(); + var second = CreateTimedEntry(12346); + Assert.True(await Insert(fake, second, 1)); + fake.Calls.Clear(); + var rowStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var releaseRow = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var activeReads = 0; + fake.BeforeRead = async path => + { + var active = Interlocked.Increment(ref activeReads); + Assert.Equal(1, active); + try + { + if (path == ZooKeeperNativeFake.RowPath(first.SiloAddress)) + { + rowStarted.TrySetResult(); + await releaseRow.Task.WaitAsync(TestContext.Current.CancellationToken); + } + } + finally + { + Interlocked.Decrement(ref activeReads); + } + }; + + var read = Read(fake); + try + { + await rowStarted.Task.WaitAsync(TestContext.Current.CancellationToken); + Assert.Equal( + new[] { "sync /", "children /", "read " + ZooKeeperNativeFake.RowPath(first.SiloAddress) }, + fake.Calls); + Assert.False(read.IsCompleted); + } + finally + { + releaseRow.TrySetResult(); + await read; + } + + var result = await read; + Assert.Equal(2, result.Version.Version); + var members = result.Members.ToDictionary(row => row.Item1.SiloAddress); + Assert.Equal(2, members.Count); + foreach (var entry in new[] { first, second }) + { + var member = members[entry.SiloAddress]; + Assert.Equal("0", member.Item2); + Assert.Equal(ZooKeeperBasedMembershipTable.Serialize(entry), ZooKeeperBasedMembershipTable.Serialize(member.Item1)); + } + Assert.Equal( + new[] { "sync /", "children /", "read " + ZooKeeperNativeFake.RowPath(first.SiloAddress), + "read " + ZooKeeperNativeFake.HeartbeatPath(first.SiloAddress), + "read " + ZooKeeperNativeFake.RowPath(second.SiloAddress), + "read " + ZooKeeperNativeFake.HeartbeatPath(second.SiloAddress), "read /" }, + fake.Calls); + } + + [Theory] + [InlineData(9, false)] + [InlineData(9, true)] + [InlineData(128, false)] + [InlineData(128, true)] + public async Task ReadAll_ReadsOneRowAtATime(int rowCount, bool cancel) + { + var fake = new ZooKeeperNativeFake(); + var entries = Enumerable.Range(0, rowCount).Select(index => CreateTimedEntry(12345 + index)).ToArray(); + for (var index = 0; index < entries.Length; index++) + { + Assert.True(await Insert(fake, entries[index], index)); + } + + fake.Calls.Clear(); + var releaseMembers = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + using var cancellation = new CancellationTokenSource(); + var firstPath = ZooKeeperNativeFake.RowPath(entries[0].SiloAddress); + fake.BeforeRead = async path => + { + if (path == firstPath) + { + await releaseMembers.Task.WaitAsync(TestContext.Current.CancellationToken); + } + }; + + var read = ZooKeeperBasedMembershipTable.ReadCoreAsync(fake.Operations, null, cancellation.Token); + var completion = Record.ExceptionAsync(() => read); + try + { + Assert.Equal( + new[] { "sync /", "children /", "read " + firstPath }, + fake.Calls); + if (cancel) + { + cancellation.Cancel(); + } + + Assert.False(read.IsCompleted); + } + finally + { + releaseMembers.TrySetResult(); + await completion; + } + + if (cancel) + { + Assert.Equal(cancellation.Token, + Assert.IsAssignableFrom(await completion).CancellationToken); + Assert.Equal(3, fake.Calls.Count); + } + else + { + Assert.Null(await completion); + var result = await read; + Assert.Equal(entries.Length, result.Version.Version); + var members = result.Members.ToDictionary(row => row.Item1.SiloAddress); + Assert.Equal(entries.Length, members.Count); + foreach (var entry in entries) + { + var member = members[entry.SiloAddress]; + Assert.Equal("0", member.Item2); + Assert.Equal(ZooKeeperBasedMembershipTable.Serialize(entry), ZooKeeperBasedMembershipTable.Serialize(member.Item1)); + } + Assert.Equal(3 + entries.Length * 2, fake.Calls.Count); + } + } + + [Fact] + public async Task Read_MissingRow_StopsBeforeLaterRowsAndPreservesFailure() + { + var (fake, first) = await CreateNativeTable(); + var second = CreateTimedEntry(12346); + Assert.True(await Insert(fake, second, 1)); + fake.Calls.Clear(); + var missingRow = new KeeperException.NoNodeException(ZooKeeperNativeFake.RowPath(first.SiloAddress)); + fake.BeforeRead = path => + { + if (path == ZooKeeperNativeFake.RowPath(first.SiloAddress)) + { + throw missingRow; + } + + return Task.CompletedTask; + }; + + Assert.Same(missingRow, await Record.ExceptionAsync(() => Read(fake))); + Assert.Equal( + new[] { "sync /", "children /", "read " + ZooKeeperNativeFake.RowPath(first.SiloAddress), "read /" }, + fake.Calls); + Assert.DoesNotContain("read " + ZooKeeperNativeFake.RowPath(second.SiloAddress), fake.Calls); + } + + [Fact] + public async Task Read_CanceledMember_PropagatesWithoutStartingHeartbeat() + { + var fake = new ZooKeeperNativeFake(); + fake.Nodes["/"] = new([], 17); + var address = CreateSiloAddress(); + using var nativeCancellation = new CancellationTokenSource(); + nativeCancellation.Cancel(); + fake.BeforeRead = path => path == ZooKeeperNativeFake.RowPath(address) + ? Task.FromCanceled(nativeCancellation.Token) + : Task.CompletedTask; + + var failure = await Assert.ThrowsAnyAsync(() => Read(fake, address)); + + Assert.Equal(nativeCancellation.Token, failure.CancellationToken); + Assert.Equal( + new[] { "sync /", "read /", "read " + ZooKeeperNativeFake.RowPath(address) }, + fake.Calls); + } + + [Fact] + public async Task Read_CancellationBeforeHeartbeat_AwaitsTheStartedMemberRequest() { var (fake, entry) = await CreateNativeTable(); + using var cancellation = new CancellationTokenSource(); + var releaseMember = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + fake.BeforeRead = async path => + { + if (path == ZooKeeperNativeFake.RowPath(entry.SiloAddress)) + { + cancellation.Cancel(); + await releaseMember.Task.WaitAsync(TestContext.Current.CancellationToken); + } + }; + + var read = ZooKeeperBasedMembershipTable.ReadCoreAsync(fake.Operations, entry.SiloAddress, cancellation.Token); + var completion = Record.ExceptionAsync(() => read); + try + { + Assert.False(read.IsCompleted); + Assert.Equal( + new[] { "sync /", "read /", "read " + ZooKeeperNativeFake.RowPath(entry.SiloAddress) }, + fake.Calls); + } + finally + { + releaseMember.TrySetResult(); + await completion; + } + + Assert.Equal(cancellation.Token, + Assert.IsAssignableFrom(await completion).CancellationToken); + } + + [Theory] + [InlineData(2)] + [InlineData(128)] + public async Task ReadAll_UpdateDuringLaterRow_RetriesTheWholeSnapshot(int rowCount) + { + var (fake, first) = await CreateNativeTable(); + var others = Enumerable.Range(1, rowCount - 1).Select(index => CreateTimedEntry(12345 + index)).ToArray(); + for (var index = 0; index < others.Length; index++) + { + Assert.True(await Insert(fake, others[index], index + 1)); + } + + var last = others[^1]; + fake.Calls.Clear(); + fake.AfterRead = async path => + { + if (path == ZooKeeperNativeFake.RowPath(last.SiloAddress)) + { + fake.AfterRead = null; + first.Status = SiloStatus.Dead; + Assert.True(await ZooKeeperBasedMembershipTable.UpdateRowCoreAsync( + fake.Operations, first, "0", + new TableVersion(rowCount + 1, rowCount.ToString(CultureInfo.InvariantCulture)), + TestContext.Current.CancellationToken)); + } + }; + + var result = await Read(fake); + + Assert.Equal(rowCount + 1, result.Version.Version); + Assert.Equal((rowCount + 1).ToString(CultureInfo.InvariantCulture), result.Version.VersionEtag); + Assert.Equal(SiloStatus.Dead, result.TryGet(first.SiloAddress)!.Item1.Status); + Assert.All(others, entry => Assert.Equal(SiloStatus.Active, result.TryGet(entry.SiloAddress)!.Item1.Status)); + Assert.Equal(rowCount, result.Members.Count); + Assert.Equal(1, fake.Calls.Count(call => call == "sync /")); + Assert.Equal(2, fake.Calls.Count(call => call == "children /")); + Assert.Equal(2, fake.Calls.Count(call => call == "read " + ZooKeeperNativeFake.RowPath(first.SiloAddress))); + Assert.All(others, entry => Assert.Equal(2, + fake.Calls.Count(call => call == "read " + ZooKeeperNativeFake.RowPath(entry.SiloAddress)))); + } + + [Fact] + public async Task ReadAll_PerpetualCanonicalChurn_StopsAfterMaximumAttempts() + { + var (fake, entry) = await CreateNativeTable(); + fake.BeforeRead = path => + { + if (path == "/") + { + var root = fake.Nodes["/"]; + fake.Nodes["/"] = root with { Version = root.Version + 1 }; + } + + return Task.CompletedTask; + }; + + var failure = await Assert.ThrowsAsync(() => Read(fake)); + + Assert.Equal( + $"Unable to read a consistent ZooKeeper membership snapshot after {ZooKeeperBasedMembershipTable.MAX_MEMBERSHIP_SNAPSHOT_ATTEMPTS} attempts.", + failure.Message); + Assert.Equal(1, fake.Calls.Count(call => call == "sync /")); + Assert.Equal(ZooKeeperBasedMembershipTable.MAX_MEMBERSHIP_SNAPSHOT_ATTEMPTS, + fake.Calls.Count(call => call == "children /")); + Assert.Equal(ZooKeeperBasedMembershipTable.MAX_MEMBERSHIP_SNAPSHOT_ATTEMPTS, + fake.Calls.Count(call => call == "read " + ZooKeeperNativeFake.RowPath(entry.SiloAddress))); + Assert.Equal(ZooKeeperBasedMembershipTable.MAX_MEMBERSHIP_SNAPSHOT_ATTEMPTS, + fake.Calls.Count(call => call == "read " + ZooKeeperNativeFake.HeartbeatPath(entry.SiloAddress))); + Assert.Equal(ZooKeeperBasedMembershipTable.MAX_MEMBERSHIP_SNAPSHOT_ATTEMPTS, + fake.Calls.Count(call => call == "read /")); + } + + [Fact] + public async Task ReadAll_StabilizesOnFinalAllowedAttempt() + { + var (fake, entry) = await CreateNativeTable(); + var closingReads = 0; + fake.BeforeRead = path => + { + if (path == "/" && ++closingReads < ZooKeeperBasedMembershipTable.MAX_MEMBERSHIP_SNAPSHOT_ATTEMPTS) + { + var root = fake.Nodes["/"]; + fake.Nodes["/"] = root with { Version = root.Version + 1 }; + } + + return Task.CompletedTask; + }; + + var result = await Read(fake); + + Assert.Single(result.Members); + Assert.Equal(entry.SiloAddress, result.Members[0].Item1.SiloAddress); + Assert.Equal(ZooKeeperBasedMembershipTable.MAX_MEMBERSHIP_SNAPSHOT_ATTEMPTS, closingReads); + Assert.Equal(ZooKeeperBasedMembershipTable.MAX_MEMBERSHIP_SNAPSHOT_ATTEMPTS, + fake.Calls.Count(call => call == "children /")); + } + + [Fact] + public async Task ReadAll_CancellationStopsBeforeRequestingTheNextRow() + { + var (fake, first) = await CreateNativeTable(); var second = CreateTimedEntry(12346); Assert.True(await Insert(fake, second, 1)); - var failure = new KeeperException.NoAuthException(); + fake.Calls.Clear(); + using var cancellation = new CancellationTokenSource(); + fake.AfterRead = path => + { + if (path == ZooKeeperNativeFake.HeartbeatPath(first.SiloAddress)) + { + cancellation.Cancel(); + } + + return Task.CompletedTask; + }; + + var exception = await Assert.ThrowsAnyAsync(() => + ZooKeeperBasedMembershipTable.ReadCoreAsync(fake.Operations, null, cancellation.Token)); + + Assert.Equal(cancellation.Token, exception.CancellationToken); + Assert.Equal( + new[] { "sync /", "children /", "read " + ZooKeeperNativeFake.RowPath(first.SiloAddress), + "read " + ZooKeeperNativeFake.HeartbeatPath(first.SiloAddress) }, + fake.Calls); + } + + [Fact] + public async Task Read_MissingRow_DoesNotObserveLaterAuthorizationFailure() + { + var (fake, entry) = await CreateNativeTable(); + var second = CreateTimedEntry(12346); + Assert.True(await Insert(fake, second, 1)); + var missing = new KeeperException.NoNodeException(ZooKeeperNativeFake.RowPath(entry.SiloAddress)); + var authorization = new KeeperException.NoAuthException(); fake.BeforeRead = path => path == ZooKeeperNativeFake.RowPath(entry.SiloAddress) - ? Task.FromException(new KeeperException.NoNodeException(path)) + ? Task.FromException(missing) : path == ZooKeeperNativeFake.RowPath(second.SiloAddress) - ? Task.FromException(failure) : Task.CompletedTask; + ? Task.FromException(authorization) : Task.CompletedTask; var actual = await Record.ExceptionAsync(() => Read(fake)); - Assert.Same(failure, actual); + Assert.Same(missing, actual); + Assert.DoesNotContain("read " + ZooKeeperNativeFake.RowPath(second.SiloAddress), fake.Calls); Assert.Equal(1, fake.Calls.Count(call => call == "sync /")); } @@ -1015,13 +1377,13 @@ await ZooKeeperBasedMembershipTable.CleanupCoreAsync( } [Fact] - public async Task Cleanup_ParallelFailure_WaitsForOutstandingRequests() + public async Task Cleanup_ReadsRowsSequentiallyAndStopsOnFailure() { var (fake, first) = await CreateNativeTable(); var second = CreateTimedEntry(12346); Assert.True(await Insert(fake, second, 1)); + fake.Calls.Clear(); var firstStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); - var failed = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); var completeFirst = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); var failure = new KeeperException.NoAuthException(); fake.BeforeRead = path => @@ -1034,7 +1396,6 @@ public async Task Cleanup_ParallelFailure_WaitsForOutstandingRequests() if (path == ZooKeeperNativeFake.RowPath(second.SiloAddress)) { - failed.SetResult(); return Task.FromException(failure); } @@ -1045,8 +1406,9 @@ public async Task Cleanup_ParallelFailure_WaitsForOutstandingRequests() fake.Operations, DateTime.UnixEpoch.AddDays(2), TestContext.Current.CancellationToken); try { - await Task.WhenAll(firstStarted.Task, failed.Task).WaitAsync(TestContext.Current.CancellationToken); + await firstStarted.Task.WaitAsync(TestContext.Current.CancellationToken); Assert.False(operation.IsCompleted); + Assert.DoesNotContain("read " + ZooKeeperNativeFake.RowPath(second.SiloAddress), fake.Calls); } finally { @@ -1054,6 +1416,10 @@ public async Task Cleanup_ParallelFailure_WaitsForOutstandingRequests() } Assert.Same(failure, await Record.ExceptionAsync(() => operation)); + Assert.Equal( + new[] { "children /", "read " + ZooKeeperNativeFake.RowPath(first.SiloAddress), + "read " + ZooKeeperNativeFake.RowPath(second.SiloAddress) }, + fake.Calls); Assert.Equal(5, fake.Nodes.Count); } @@ -1146,6 +1512,61 @@ public async Task MembershipOperations_InfrastructureFailure_PropagatesSameExcep Assert.Equal(1, fake.Nodes["/"].Version); } + [Fact] + public async Task Cleanup_PerpetualConflict_StopsAfterMaximumAttempts() + { + var (fake, entry) = await CreateNativeTable(SiloStatus.Dead); + var later = CreateTimedEntry(12346); + Assert.True(await Insert(fake, later, 1)); + fake.Calls.Clear(); + fake.Transactions.Clear(); + fake.BeforeMulti = _ => Task.FromException( + new KeeperException.BadVersionException(ZooKeeperNativeFake.HeartbeatPath(entry.SiloAddress))); + + var failure = await Assert.ThrowsAsync(() => + ZooKeeperBasedMembershipTable.CleanupCoreAsync( + fake.Operations, + DateTime.UnixEpoch.AddDays(2), + TestContext.Current.CancellationToken)); + + Assert.Equal( + $"Unable to clean ZooKeeper membership row '{ZooKeeperNativeFake.RowPath(entry.SiloAddress)}' after {ZooKeeperBasedMembershipTable.MAX_CLEANUP_ROW_ATTEMPTS} concurrent modifications.", + failure.Message); + Assert.Equal(ZooKeeperBasedMembershipTable.MAX_CLEANUP_ROW_ATTEMPTS, fake.Transactions.Count); + Assert.Equal(ZooKeeperBasedMembershipTable.MAX_CLEANUP_ROW_ATTEMPTS, + fake.Calls.Count(call => call == "read " + ZooKeeperNativeFake.RowPath(entry.SiloAddress))); + Assert.Equal(ZooKeeperBasedMembershipTable.MAX_CLEANUP_ROW_ATTEMPTS, + fake.Calls.Count(call => call == "read " + ZooKeeperNativeFake.HeartbeatPath(entry.SiloAddress))); + Assert.True(fake.Nodes.ContainsKey(ZooKeeperNativeFake.RowPath(entry.SiloAddress))); + Assert.True(fake.Nodes.ContainsKey(ZooKeeperNativeFake.HeartbeatPath(entry.SiloAddress))); + Assert.DoesNotContain("read " + ZooKeeperNativeFake.RowPath(later.SiloAddress), fake.Calls); + } + + [Fact] + public async Task Cleanup_ConflictOnFinalAllowedAttempt_ThenSucceeds() + { + var (fake, entry) = await CreateNativeTable(SiloStatus.Dead); + var attempts = 0; + fake.BeforeMulti = _ => + { + if (++attempts < ZooKeeperBasedMembershipTable.MAX_CLEANUP_ROW_ATTEMPTS) + { + throw new KeeperException.BadVersionException(ZooKeeperNativeFake.HeartbeatPath(entry.SiloAddress)); + } + + return Task.CompletedTask; + }; + + await ZooKeeperBasedMembershipTable.CleanupCoreAsync( + fake.Operations, + DateTime.UnixEpoch.AddDays(2), + TestContext.Current.CancellationToken); + + Assert.Equal(ZooKeeperBasedMembershipTable.MAX_CLEANUP_ROW_ATTEMPTS, attempts); + Assert.Single(fake.Nodes); + Assert.Equal("/", Assert.Single(fake.Nodes).Key); + } + [Theory] [InlineData("update")] [InlineData("cleanup")] diff --git a/test/Extensions/Orleans.Clustering.ZooKeeper.Tests/ZooKeeperNativeDiagnostics.cs b/test/Extensions/Orleans.Clustering.ZooKeeper.Tests/ZooKeeperNativeDiagnostics.cs new file mode 100644 index 00000000000..692944b9b80 --- /dev/null +++ b/test/Extensions/Orleans.Clustering.ZooKeeper.Tests/ZooKeeperNativeDiagnostics.cs @@ -0,0 +1,98 @@ +using System.Collections.Concurrent; +using System.Diagnostics; +using System.Diagnostics.Tracing; +using org.apache.utils; +using org.apache.zookeeper; + +namespace UnitTests.MembershipTests; + +internal sealed class NativeSocketDiagnostics(Action write) : EventListener +{ + private const int MaxLoggedEvents = 256; + private readonly Action _write = write; + private int _eventCount; + + protected override void OnEventSourceCreated(EventSource eventSource) + { + if (eventSource.Name == "Private.InternalDiagnostics.System.Net.Sockets") + { + // Capture native completion errors before the SDK assigns its own SocketError. + EnableEvents(eventSource, EventLevel.Error, (EventKeywords)1); + _write($"{DateTime.UtcNow:O} [NativeSockets#{GetHashCode()}] Error listener enabled"); + } + } + + protected override void OnEventWritten(EventWrittenEventArgs eventData) + { + if (eventData.EventName is "ErrorMessage" or "EventSourceMessage") + { + var count = Interlocked.Increment(ref _eventCount); + if (count <= MaxLoggedEvents) + { + _write($"{DateTime.UtcNow:O} [NativeSockets#{GetHashCode()}] {eventData.EventName}: {string.Join("; ", eventData.Payload!)}"); + } + else if (count == MaxLoggedEvents + 1) + { + _write($"{DateTime.UtcNow:O} [NativeSockets#{GetHashCode()}] Event limit reached; subsequent event details are omitted"); + } + } + } + + public override void Dispose() + { + base.Dispose(); + _write($"{DateTime.UtcNow:O} [NativeSockets#{GetHashCode()}] Listener disposed; observed-events={Volatile.Read(ref _eventCount)}; detail-limit={MaxLoggedEvents}"); + } +} + +internal sealed class ZooKeeperNativeDiagnostics : ILogConsumer, IDisposable +{ + private const int MaxSdkMessages = 256; + private readonly TraceLevel _previousLevel = ZooKeeper.LogLevel; + private readonly bool _previousTrace = ZooKeeper.LogToTrace; + private readonly ILogConsumer? _previousConsumer = ZooKeeper.CustomLogConsumer; + private readonly ConcurrentQueue _messages = new(); + private readonly NativeSocketDiagnostics _sockets; + private int _sdkMessageCount; + + internal ZooKeeperNativeDiagnostics() + { + _sockets = new NativeSocketDiagnostics(_messages.Enqueue); + ZooKeeper.LogLevel = TraceLevel.Info; + ZooKeeper.LogToTrace = false; + ZooKeeper.CustomLogConsumer = this; + } + + public void Log(TraceLevel severity, string className, string message, Exception exception) + { + if (exception is not null || severity <= TraceLevel.Warning) + { + var record = $"{DateTime.UtcNow:O} [{className}] {severity}: {message}{Environment.NewLine}{exception}"; + Console.Error.WriteLine(record); + var count = Interlocked.Increment(ref _sdkMessageCount); + if (count <= MaxSdkMessages) + { + _messages.Enqueue(record); + } + else if (count == MaxSdkMessages + 1) + { + _messages.Enqueue($"{DateTime.UtcNow:O} SDK message limit reached; subsequent message details are omitted"); + } + } + } + + public void Dispose() + { + ZooKeeper.CustomLogConsumer = _previousConsumer; + ZooKeeper.LogToTrace = _previousTrace; + ZooKeeper.LogLevel = _previousLevel; + _sockets.Dispose(); + _messages.Enqueue($"SDK messages observed={Volatile.Read(ref _sdkMessageCount)}; detail-limit={MaxSdkMessages}"); + } + + internal async Task WriteAsync(string path) + { + Directory.CreateDirectory(Path.GetDirectoryName(path)!); + await File.WriteAllLinesAsync(path, _messages); + } +} diff --git a/test/Extensions/Orleans.Clustering.ZooKeeper.Tests/ZooKeeperReadResilienceTests.cs b/test/Extensions/Orleans.Clustering.ZooKeeper.Tests/ZooKeeperReadResilienceTests.cs new file mode 100644 index 00000000000..18928c70fa3 --- /dev/null +++ b/test/Extensions/Orleans.Clustering.ZooKeeper.Tests/ZooKeeperReadResilienceTests.cs @@ -0,0 +1,189 @@ +using System.Collections.Concurrent; +using System.Diagnostics; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Options; +using Orleans.Clustering.TestKit; +using Orleans.Configuration; +using Orleans.Runtime.Membership; +using Orleans.TestingHost.Utils; +using org.apache.zookeeper; +using TestExtensions; +using Xunit; + +namespace UnitTests.MembershipTests; + +[Collection(TestEnvironmentFixture.DefaultCollection)] +[TestCategory("Membership"), TestCategory("ZooKeeper")] +[TestSuite("Functional"), TestProvider("ZooKeeper"), TestArea("Membership")] +public sealed class ZooKeeperReadResilienceTests : IAsyncLifetime +{ + private readonly ZooKeeperNativeDiagnostics _diagnostics = new(); + private readonly ConcurrentBag _sessions = []; + private readonly List _fixtures = []; + private readonly ConcurrentBag _legacyOperations = []; + private readonly ConcurrentBag _probes = []; + private readonly string _socketLog = Path.Combine(AppContext.BaseDirectory, "TestResults", $"zookeeper-sockets-{Guid.NewGuid():N}.log"); + private readonly ILoggerFactory _loggerFactory = TestingUtils.CreateDefaultLoggerFactory( + TestingUtils.CreateTraceFileName("zookeeper-reads", Guid.NewGuid().ToString("N")), new LoggerFilterOptions()); + private string _connectionString = null!; + + public async ValueTask InitializeAsync() + { + Assert.False(string.IsNullOrWhiteSpace(TestDefaultConfiguration.ZooKeeperConnectionString), + "ZooKeeper resilience tests require a configured connection string."); + _connectionString = TestDefaultConfiguration.ZooKeeperConnectionString!; + var probe = ZooKeeper.Using(_connectionString, 2000, new ConformanceWatcher(), + async client => await client.existsAsync("/", false) is not null); + _probes.Add(probe); + Assert.True(await probe.WaitAsync(TestContext.Current.CancellationToken), + "ZooKeeper resilience tests require the configured ZooKeeper service."); + } + + public async ValueTask DisposeAsync() + { + try + { + // Fixture teardown and canceled caller waits can return before native close. + await Task.WhenAll(_fixtures.Select(fixture => DrainFixtureAsync(fixture.DisposeAsync)) + .Concat(_sessions.Select(session => session.Completion)) + .Concat(_legacyOperations) + .Concat(_probes)); + } + finally + { + _diagnostics.Dispose(); + await _diagnostics.WriteAsync(_socketLog); + _loggerFactory.Dispose(); + } + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task MembershipTable_ZooKeeper_RepeatedSnapshotReadCompatibility(bool pointRead) + { + for (var iteration = 0; iteration < 3; iteration++) + { + var started = Stopwatch.GetTimestamp(); + await CreateFixture().RunAsync(async (fixture, cancellationToken) => + { + var phase = "concurrent scenario"; + await CapturePrimaryScenarioFailureAsync(async () => + { + var runner = new MembershipTableTestRunner(fixture, seed: 17, concurrencyRowCount: 128); + if (pointRead) + { + await runner.ConcurrentReadRow_ReturnsOnlyAtomicCommittedViews(cancellationToken); + } + else + { + await runner.ConcurrentReadAll_ReturnsOnlyAtomicCommittedViews(cancellationToken); + } + + phase = "stable ReadAll"; + var readStarted = Stopwatch.GetTimestamp(); + var snapshot = await fixture.First.ReadAllAsync(cancellationToken); + TestContext.Current.TestOutputHelper?.WriteLine( + $"Stable snapshot rows={snapshot.Members.Count}; elapsed={Stopwatch.GetElapsedTime(readStarted)}"); + }, failure => RecordPrimaryFailure( + $"Primary snapshot failure; pointRead={pointRead}; repetition={iteration + 1}/3; phase={phase}", failure)); + }, TestContext.Current.CancellationToken); + TestContext.Current.TestOutputHelper?.WriteLine( + $"Snapshot compatibility pass {iteration + 1}/3; pointRead={pointRead}; elapsed={Stopwatch.GetElapsedTime(started)}"); + } + } + + [Fact] + public Task MembershipTable_ZooKeeper_DeletionProbe_PreservesOriginalScope() => + CreateFixture().RunAsync((fixture, cancellationToken) => + CapturePrimaryScenarioFailureAsync( + () => new MembershipTableTestRunner(fixture, seed: 17) + .DeleteMembershipTableEntries_DeletesOwnClusterAndPreservesOtherCluster(cancellationToken), + failure => RecordPrimaryFailure("Primary deletion-probe failure", failure)), + TestContext.Current.CancellationToken); + + private void RecordPrimaryFailure(string context, string failure) + { + var record = context + Environment.NewLine + failure; + _loggerFactory.CreateLogger().LogError("{PrimaryFailure}", record); + TestContext.Current.TestOutputHelper?.WriteLine(record); + } + + internal static async Task CapturePrimaryScenarioFailureAsync(Func action, Action record) + { + try + { + await action(); + } + catch (Exception exception) + { + // Keep the caller stack before teardown re-observes retained operation failures. + record(exception.ToString()); + throw; + } + } + + internal static async Task DrainFixtureAsync(Func dispose) + { + try + { + await dispose(); + } + catch (TimeoutException exception) when (exception.Data["ClusteringTestKit.CleanupCompletion"] is Task completion) + { + await completion; + throw; + } + } + + private MembershipTableTestFixture CreateFixture() + { + var fixture = new MembershipTableTestFixture(nameof(ZooKeeperReadResilienceTests), (serviceId, clusterId, cancellationToken) => + { + cancellationToken.ThrowIfCancellationRequested(); + var logger = _loggerFactory.CreateLogger(); + var sessions = new ConcurrentBag(); + var legacyOperations = new ConcurrentBag(); + var table = new ZooKeeperBasedMembershipTable( + logger, + Options.Create(new ZooKeeperClusteringSiloOptions { ConnectionString = _connectionString }), + Options.Create(new ClusterOptions { ServiceId = serviceId, ClusterId = clusterId }), + readOnly => + { + var session = ZooKeeperBasedMembershipTable.CreateSession( + _connectionString + "/" + clusterId, new ZooKeeperWatcher(logger), readOnly); + sessions.Add(session); + _sessions.Add(session); + return session; + }, + null, + operation => + { + legacyOperations.Add(operation); + _legacyOperations.Add(operation); + }); + return ValueTask.FromResult(new MembershipTableTestHandle(table, + () => new ValueTask(Task.WhenAll(sessions.Select(session => session.Completion).Concat(legacyOperations))))); + }, IsConformanceClusterDeletedAsync); + _fixtures.Add(fixture); + return fixture; + } + + private async ValueTask IsConformanceClusterDeletedAsync(string clusterId, CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + var probe = ZooKeeper.Using(_connectionString, 10_000, new ConformanceWatcher(), async client => + { + await client.sync("/"); + cancellationToken.ThrowIfCancellationRequested(); + return await client.existsAsync("/" + clusterId, false) is null; + }); + _probes.Add(probe); + return await probe; + } + + private sealed class ConformanceWatcher : Watcher + { + public override Task process(WatchedEvent @event) => Task.CompletedTask; + } +} diff --git a/test/Extensions/Orleans.Clustering.ZooKeeper.Tests/ZooKeeperReadRetryTests.cs b/test/Extensions/Orleans.Clustering.ZooKeeper.Tests/ZooKeeperReadRetryTests.cs new file mode 100644 index 00000000000..c48c978f5c0 --- /dev/null +++ b/test/Extensions/Orleans.Clustering.ZooKeeper.Tests/ZooKeeperReadRetryTests.cs @@ -0,0 +1,1061 @@ +using System.Collections.Concurrent; +using System.Globalization; +using System.Net; +using System.Threading.Channels; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Logging.Abstractions; +using Microsoft.Extensions.Options; +using Microsoft.Extensions.Time.Testing; +using Orleans.Configuration; +using Orleans.Runtime; +using Orleans.Runtime.Membership; +using org.apache.zookeeper; +using Polly; +using TestExtensions; +using Xunit; + +namespace UnitTests.MembershipTests; + +[TestCategory("Membership"), TestCategory("ZooKeeper")] +[TestSuite("BVT"), TestProvider("ZooKeeper"), TestArea("Membership")] +public sealed class ZooKeeperReadRetryTests +{ + [Theory] + [InlineData("sync", false)] + [InlineData("children", false)] + [InlineData("before", true)] + [InlineData("member", false)] + [InlineData("heartbeat", false)] + [InlineData("after", false)] + public async Task Read_ConnectionLossThenSuccess_RetriesOnlyFailedNativeRequest(string boundary, bool point) + { + var harness = await Harness.CreateAsync(); + var key = boundary switch + { + "sync" => "Sync /", + "children" => "GetChildren /", + "before" or "after" => "GetData /", + "member" => "GetData " + ZooKeeperNativeFake.RowPath(harness.Entries[0].SiloAddress), + "heartbeat" => "GetData " + ZooKeeperNativeFake.HeartbeatPath(harness.Entries[0].SiloAddress), + _ => throw new ArgumentOutOfRangeException(nameof(boundary)) + }; + var failure = new KeeperException.ConnectionLossException(); + var failed = false; + harness.BeforeRequest = request => + { + if (request == key && !failed) + { + failed = true; + throw failure; + } + return Task.CompletedTask; + }; + + var read = harness.Read(point, TestContext.Current.CancellationToken); + await harness.Clock.AdvanceNextAsync(TimeSpan.FromMilliseconds(250)); + var result = await read; + + harness.AssertSnapshot(result, point ? [harness.Entries[0]] : harness.Entries, 2); + var expected = harness.ExpectedReadCalls(point).GroupBy(value => value).ToDictionary(group => group.Key, group => group.Count()); + expected[key]++; + Assert.Equal(expected.OrderBy(pair => pair.Key), harness.CountCalls().OrderBy(pair => pair.Key)); + var warning = Assert.Single(harness.Logger.Warnings); + Assert.Same(failure, warning.Exception); + Assert.Equal(key.Split(' ')[0], warning.Values["Operation"]); + Assert.Equal(1, warning.Values["Retry"]); + Assert.Equal(250d, warning.Values["DelayMilliseconds"]); + harness.AssertOneOwner(readOnly: true); + } + + [Fact] + public async Task ReadRetry_Exhaustion_PreservesFinalExceptionAndBackoff() + { + var harness = await Harness.CreateAsync(); + var failures = Enumerable.Range(0, 5).Select(_ => new KeeperException.ConnectionLossException()).ToArray(); + var times = new List(); + harness.BeforeRequest = _ => + { + times.Add(harness.Clock.GetUtcNow()); + throw failures[times.Count - 1]; + }; + var read = harness.Read(TestContext.Current.CancellationToken); + var completion = Record.ExceptionAsync(() => read); + foreach (var delay in new[] { 250, 500, 1000, 2000 }) + { + var timer = await harness.Clock.NextTimerAsync(); + Assert.Equal(TimeSpan.FromMilliseconds(delay), timer); + var attempts = harness.Calls.Count; + harness.Clock.Advance(timer - TimeSpan.FromMilliseconds(1)); + Assert.Equal(attempts, harness.Calls.Count); + harness.Clock.Advance(TimeSpan.FromMilliseconds(1)); + } + + Assert.Same(failures[^1], await completion); + Assert.Equal(new[] { 0d, 250d, 750d, 1750d, 3750d }, times.Select(time => (time - times[0]).TotalMilliseconds)); + Assert.Equal(Enumerable.Repeat("Sync /", 5), harness.Calls); + Assert.Equal(failures.Take(4), harness.Logger.Warnings.Select(warning => warning.Exception)); + Assert.Equal(new[] { 1, 2, 3, 4 }, harness.Logger.Warnings.Select(warning => (int)warning.Values["Retry"]!)); + Assert.All(harness.Logger.Warnings, warning => Assert.Equal(4, warning.Values["MaxRetries"])); + harness.AssertOneOwner(readOnly: true); + } + + [Theory] + [InlineData("authorization")] + [InlineData("session")] + [InlineData("missing")] + [InlineData("version")] + [InlineData("cancellation")] + [InlineData("ordinary")] + public async Task ReadRetry_IneligibleFailure_PropagatesAfterOneAttempt(string kind) + { + var harness = await Harness.CreateAsync(); + Exception failure = kind switch + { + "authorization" => new KeeperException.NoAuthException(), + "session" => new KeeperException.SessionExpiredException(), + "missing" => new KeeperException.NoNodeException("/"), + "version" => new KeeperException.BadVersionException("/"), + "cancellation" => new OperationCanceledException(), + "ordinary" => new InvalidOperationException("native failure"), + _ => throw new ArgumentOutOfRangeException(nameof(kind)) + }; + harness.BeforeRequest = _ => Task.FromException(failure); + + Assert.Same(failure, await Record.ExceptionAsync(() => harness.Read(TestContext.Current.CancellationToken))); + Assert.Equal("Sync /", Assert.Single(harness.Calls)); + Assert.Empty(harness.Logger.Warnings); + harness.AssertOneOwner(readOnly: true); + } + + [Fact] + public async Task ConnectionMonitor_FailureBeforeDisconnected_WaitsForNextConnection() + { + var monitor = CreateConnectionMonitor(); + await Signal(monitor, Watcher.Event.KeeperState.SyncConnected); + var connectedGeneration = ((IZooKeeperConnectionMonitor)monitor).CaptureAttemptGeneration(); + + var wait = ((IZooKeeperConnectionMonitor)monitor) + .WaitForConnectionAfterAsync(connectedGeneration, TestContext.Current.CancellationToken).AsTask(); + Assert.False(wait.IsCompleted); + + await Signal(monitor, Watcher.Event.KeeperState.Disconnected); + Assert.False(wait.IsCompleted); + + await Signal(monitor, Watcher.Event.KeeperState.SyncConnected); + Assert.True(await wait); + } + + [Fact] + public async Task ConnectionMonitor_ReconnectBeforeWaiterRegistration_CompletesImmediately() + { + var monitor = CreateConnectionMonitor(); + await Signal(monitor, Watcher.Event.KeeperState.SyncConnected); + var connectedGeneration = ((IZooKeeperConnectionMonitor)monitor).CaptureAttemptGeneration(); + await Signal(monitor, Watcher.Event.KeeperState.Disconnected); + await Signal(monitor, Watcher.Event.KeeperState.SyncConnected); + + var wait = ((IZooKeeperConnectionMonitor)monitor) + .WaitForConnectionAfterAsync(connectedGeneration, TestContext.Current.CancellationToken); + + Assert.True(wait.IsCompletedSuccessfully); + Assert.True(await wait); + } + + [Fact] + public async Task ConnectionMonitor_OneReconnectReleasesAllWaiters_AndCanceledWaiterDoesNotPoisonPeers() + { + var monitor = CreateConnectionMonitor(); + await Signal(monitor, Watcher.Event.KeeperState.SyncConnected); + var connectedGeneration = ((IZooKeeperConnectionMonitor)monitor).CaptureAttemptGeneration(); + using var canceled = CancellationTokenSource.CreateLinkedTokenSource(TestContext.Current.CancellationToken); + var canceledWaiter = ((IZooKeeperConnectionMonitor)monitor) + .WaitForConnectionAfterAsync(connectedGeneration, canceled.Token).AsTask(); + var peers = Enumerable.Range(0, 16).Select(_ => ((IZooKeeperConnectionMonitor)monitor) + .WaitForConnectionAfterAsync(connectedGeneration, TestContext.Current.CancellationToken).AsTask()).ToArray(); + + canceled.Cancel(); + var cancellation = await Assert.ThrowsAnyAsync(() => canceledWaiter); + Assert.Equal(canceled.Token, cancellation.CancellationToken); + Assert.All(peers, peer => Assert.False(peer.IsCompleted)); + + await Signal(monitor, Watcher.Event.KeeperState.Disconnected); + Assert.All(peers, peer => Assert.False(peer.IsCompleted)); + await Signal(monitor, Watcher.Event.KeeperState.SyncConnected); + + Assert.All(await Task.WhenAll(peers), Assert.True); + } + + [Fact] + public async Task ConnectionMonitor_TimeoutReturnsWithoutReplacingNativeFailure() + { + var timeProvider = new FakeTimeProvider(); + var monitor = CreateConnectionMonitor(timeProvider, TimeSpan.FromSeconds(1)); + await Signal(monitor, Watcher.Event.KeeperState.SyncConnected); + var connectedGeneration = ((IZooKeeperConnectionMonitor)monitor).CaptureAttemptGeneration(); + var wait = ((IZooKeeperConnectionMonitor)monitor) + .WaitForConnectionAfterAsync(connectedGeneration, TestContext.Current.CancellationToken).AsTask(); + + timeProvider.Advance(TimeSpan.FromMilliseconds(999)); + Assert.False(wait.IsCompleted); + timeProvider.Advance(TimeSpan.FromMilliseconds(1)); + + Assert.False(await wait); + } + + [Fact] + public async Task ConnectionMonitor_InitialConnectionDoesNotSatisfyPostFailureReconnect() + { + var monitor = CreateConnectionMonitor(); + var attemptGeneration = ((IZooKeeperConnectionMonitor)monitor).CaptureAttemptGeneration(); + await Signal(monitor, Watcher.Event.KeeperState.SyncConnected); + + var wait = ((IZooKeeperConnectionMonitor)monitor) + .WaitForConnectionAfterAsync(attemptGeneration, TestContext.Current.CancellationToken).AsTask(); + + Assert.False(wait.IsCompleted); + await Signal(monitor, Watcher.Event.KeeperState.Disconnected); + Assert.False(wait.IsCompleted); + await Signal(monitor, Watcher.Event.KeeperState.SyncConnected); + Assert.True(await wait); + } + + [Fact] + public async Task ReadRetry_ReconnectTimeout_ReissuesAndPreservesNativeFailure() + { + var harness = await Harness.CreateAsync(); + var monitorTime = new FakeTimeProvider(); + var monitor = CreateConnectionMonitor(monitorTime, TimeSpan.FromSeconds(1)); + harness.ConnectionMonitor = monitor; + await Signal(monitor, Watcher.Event.KeeperState.SyncConnected); + var failure = new KeeperException.ConnectionLossException(); + harness.BeforeRequest = _ => Task.FromException(failure); + var read = harness.Read(TestContext.Current.CancellationToken); + var completion = Record.ExceptionAsync(() => read); + + foreach (var delay in new[] { 250, 500, 1000, 2000 }) + { + await harness.Clock.AdvanceNextAsync(TimeSpan.FromMilliseconds(delay)); + monitorTime.Advance(TimeSpan.FromSeconds(1)); + } + + Assert.Same(failure, await completion); + Assert.Equal(Enumerable.Repeat("Sync /", 5), harness.Calls); + harness.AssertOneOwner(readOnly: true); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task ReadOwner_PreCanceled_DoesNotCreateSession(bool point) + { + var harness = await Harness.CreateAsync(); + using var cancellation = new CancellationTokenSource(); + cancellation.Cancel(); + + var failure = await Assert.ThrowsAnyAsync(() => harness.Read(point, cancellation.Token)); + + Assert.Equal(cancellation.Token, failure.CancellationToken); + Assert.Empty(harness.Sessions); + Assert.Empty(harness.Calls); + Assert.Equal(0, harness.CloseCount); + } + + [Fact] + public async Task ReadRetry_CancellationDuringDelay_StopsAdmission() + { + var harness = await Harness.CreateAsync(); + using var cancellation = new CancellationTokenSource(); + harness.BeforeRequest = _ => Task.FromException(new KeeperException.ConnectionLossException()); + var read = harness.Read(cancellationToken: cancellation.Token); + Assert.Equal(TimeSpan.FromMilliseconds(250), await harness.Clock.NextTimerAsync()); + + cancellation.Cancel(); + + var failure = await Assert.ThrowsAnyAsync(() => read); + Assert.Equal(cancellation.Token, failure.CancellationToken); + await Assert.ThrowsAnyAsync(() => Assert.Single(harness.Sessions).Completion); + harness.Clock.Advance(TimeSpan.FromDays(1)); + Assert.Equal("Sync /", Assert.Single(harness.Calls)); + Assert.Single(harness.Logger.Warnings); + harness.AssertOneOwner(readOnly: true); + } + + [Fact] + public async Task ReadOwner_CanceledCaller_JoinsNativeTasksBeforeClose() + { + var harness = await Harness.CreateAsync(); + using var cancellation = new CancellationTokenSource(); + var first = Gate(); + var firstCompleted = Gate(); + var closeStarted = Gate(); + var releaseClose = Gate(); + var firstKey = "GetData " + ZooKeeperNativeFake.RowPath(harness.Entries[0].SiloAddress); + var secondKey = "GetData " + ZooKeeperNativeFake.RowPath(harness.Entries[1].SiloAddress); + harness.BeforeRequest = request => request == firstKey ? first.Task : Task.CompletedTask; + harness.AfterRequest = request => + { + if (request == firstKey) + firstCompleted.SetResult(); + }; + harness.Close = () => + { + closeStarted.SetResult(); + return releaseClose.Task; + }; + + var read = harness.Read(cancellationToken: cancellation.Token); + var owner = Assert.Single(harness.Sessions); + try + { + Assert.Equal(new[] { "Sync /", "GetChildren /", firstKey }, harness.Calls); + cancellation.Cancel(); + var failure = await Assert.ThrowsAnyAsync(() => read); + Assert.Equal(cancellation.Token, failure.CancellationToken); + Assert.False(owner.Completion.IsCompleted); + Assert.False(closeStarted.Task.IsCompleted); + first.SetResult(); + await firstCompleted.Task.WaitAsync(TestContext.Current.CancellationToken); + Assert.False(owner.Completion.IsCompleted); + Assert.False(closeStarted.Task.IsCompleted); + await closeStarted.Task.WaitAsync(TestContext.Current.CancellationToken); + Assert.False(owner.Completion.IsCompleted); + Assert.DoesNotContain(secondKey, harness.Calls); + } + finally + { + first.TrySetResult(); + releaseClose.TrySetResult(); + await Record.ExceptionAsync(() => owner.Completion); + } + + var ownedFailure = await Assert.ThrowsAnyAsync(() => owner.Completion); + Assert.Equal(cancellation.Token, ownedFailure.CancellationToken); + Assert.Equal(new[] { "Sync /", "GetChildren /", firstKey }, harness.Calls); + harness.AssertOneOwner(readOnly: true); + } + + [Fact] + public async Task ReadOwner_CloseCompletesBeforeResult() + { + var harness = await Harness.CreateAsync(); + var close = Gate(); + harness.Close = () => close.Task; + var read = harness.Read(TestContext.Current.CancellationToken); + try + { + Assert.Equal(harness.ExpectedReadCalls(), harness.Calls); + Assert.Equal(1, harness.CloseCount); + Assert.False(read.IsCompleted); + Assert.False(Assert.Single(harness.Sessions).Completion.IsCompleted); + } + finally + { + close.TrySetResult(); + } + harness.AssertSnapshot(await read, harness.Entries, 2); + harness.AssertOneOwner(readOnly: true); + } + + [Fact] + public async Task ReadOwner_CloseFailureAfterSuccess_DoesNotReplay() + { + var harness = await Harness.CreateAsync(); + var failure = new KeeperException.ConnectionLossException(); + harness.Close = () => Task.FromException(failure); + + Assert.Same(failure, await Record.ExceptionAsync(() => harness.Read(TestContext.Current.CancellationToken))); + + Assert.Equal(harness.ExpectedReadCalls(), harness.Calls); + Assert.Empty(harness.Logger.Warnings); + harness.AssertOneOwner(readOnly: true); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task ReadOwner_CompletionIsPendingAtPublicationAndSynchronousPrefix(bool checkAtPublication) + { + var harness = await Harness.CreateAsync(); + var release = Gate(); + Task completionAtPublication = null!; + Task snapshotDrain = null!; + harness.OnSessionCreated = session => + { + completionAtPublication = session.Completion; + snapshotDrain = Task.WhenAll(harness.Sessions.Select(owned => owned.Completion)); + if (checkAtPublication) + { + Assert.False(completionAtPublication.IsCompleted); + Assert.False(snapshotDrain.IsCompleted); + } + }; + harness.BeforeRequest = request => + { + if (request == "Sync /") + { + Assert.Same(completionAtPublication, Assert.Single(harness.Sessions).Completion); + Assert.False(completionAtPublication.IsCompleted); + Assert.False(snapshotDrain.IsCompleted); + return release.Task; + } + return Task.CompletedTask; + }; + var read = harness.Read(TestContext.Current.CancellationToken); + try + { + Assert.False(read.IsCompleted); + Assert.False(snapshotDrain.IsCompleted); + Assert.Equal(0, harness.CloseCount); + Assert.Same(completionAtPublication, Assert.Single(harness.Sessions).Completion); + } + finally + { + release.TrySetResult(); + } + + harness.AssertSnapshot(await read, harness.Entries, 2); + await snapshotDrain.WaitAsync(TestContext.Current.CancellationToken); + Assert.Equal(harness.ExpectedReadCalls(), harness.Calls); + Assert.Same(completionAtPublication, Assert.Single(harness.Sessions).Completion); + harness.AssertOneOwner(readOnly: true); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task ReadOwner_CallbackAndCloseFailures_PreserveBothExceptions(bool point) + { + var harness = await Harness.CreateAsync(); + var primary = new KeeperException.NoAuthException(); + var secondary = new KeeperException.ConnectionLossException(); + harness.BeforeRequest = _ => Task.FromException(primary); + harness.Close = () => Task.FromException(secondary); + + var actual = await Assert.ThrowsAsync(() => harness.Read(point, TestContext.Current.CancellationToken)); + + Assert.Equal(new Exception[] { primary, secondary }, actual.InnerExceptions); + Assert.Same(actual, await Record.ExceptionAsync(() => Assert.Single(harness.Sessions).Completion)); + Assert.Equal("Sync /", Assert.Single(harness.Calls)); + Assert.Empty(harness.Logger.Warnings); + harness.AssertOneOwner(readOnly: true); + } + + [Theory] + [InlineData(false, false)] + [InlineData(false, true)] + [InlineData(true, false)] + [InlineData(true, true)] + public async Task Read_ConnectionLossDuringConcurrentMutation_RefencesWholeSnapshot(bool point, bool cleanup) + { + var harness = await Harness.CreateAsync(); + var entry = harness.Entries[0]; + var heartbeat = "GetData " + ZooKeeperNativeFake.HeartbeatPath(entry.SiloAddress); + var failed = false; + harness.BeforeRequest = request => + { + if (request == heartbeat && !failed) + { + failed = true; + throw new KeeperException.ConnectionLossException(); + } + return Task.CompletedTask; + }; + var read = harness.Read(point, TestContext.Current.CancellationToken); + var delay = await harness.Clock.NextTimerAsync(); + if (cleanup) + { + harness.Fake.Nodes.Remove(ZooKeeperNativeFake.RowPath(entry.SiloAddress)); + harness.Fake.Nodes.Remove(ZooKeeperNativeFake.HeartbeatPath(entry.SiloAddress)); + var root = harness.Fake.Nodes["/"]; + harness.Fake.Nodes["/"] = root with { ChildrenVersion = root.ChildrenVersion + 1 }; + } + else + { + entry.Status = SiloStatus.Dead; + Assert.True(await ZooKeeperBasedMembershipTable.UpdateRowCoreAsync( + harness.Fake.Operations, entry, "0", new TableVersion(3, "2"), TestContext.Current.CancellationToken)); + } + harness.Clock.Advance(delay); + + var result = await read; + var entries = cleanup ? harness.Entries.Skip(1).ToArray() : harness.Entries; + harness.AssertSnapshot(result, point ? entries.Where(value => value.SiloAddress.Equals(entry.SiloAddress)) : entries, + cleanup ? 2 : 3); + Assert.Equal(1, harness.Calls.Count(call => call == "Sync /")); + Assert.Equal(point ? 4 : 2, harness.Calls.Count(call => call == (point ? "GetData /" : "GetChildren /"))); + Assert.Equal(cleanup && !point ? 1 : 2, + harness.Calls.Count(call => call == "GetData " + ZooKeeperNativeFake.RowPath(entry.SiloAddress))); + harness.AssertOneOwner(readOnly: true); + } + + [Theory] + [InlineData(9)] + [InlineData(128)] + public async Task Read_RetryingRow_ContinuesSequentialSnapshotAfterRecovery(int rowCount) + { + var harness = await Harness.CreateAsync(rowCount); + var key = "GetData " + ZooKeeperNativeFake.RowPath(harness.Entries[0].SiloAddress); + var failed = false; + harness.BeforeRequest = request => + { + if (request == key && !failed) + { + failed = true; + throw new KeeperException.ConnectionLossException(); + } + return Task.CompletedTask; + }; + var read = harness.Read(TestContext.Current.CancellationToken); + Assert.False(read.IsCompleted); + Assert.All(harness.Entries.Skip(1), entry => + Assert.DoesNotContain("GetData " + ZooKeeperNativeFake.RowPath(entry.SiloAddress), harness.Calls)); + await harness.Clock.AdvanceNextAsync(TimeSpan.FromMilliseconds(250)); + + harness.AssertSnapshot(await read, harness.Entries, rowCount); + var expected = harness.ExpectedReadCalls().ToDictionary(call => call, _ => 1); + expected[key]++; + Assert.Equal(expected.OrderBy(pair => pair.Key), harness.CountCalls().OrderBy(pair => pair.Key)); + harness.AssertOneOwner(readOnly: true); + } + + [Theory] + [InlineData(9, false)] + [InlineData(9, true)] + [InlineData(128, false)] + [InlineData(128, true)] + public async Task ReadOwner_ReadsOneRowAtATimeAndOwnsAdmittedRequest(int rowCount, bool cancel) + { + var harness = await Harness.CreateAsync(rowCount); + using var cancellation = new CancellationTokenSource(); + var release = Gate(); + harness.BeforeRequest = request => request.StartsWith("GetData ", StringComparison.Ordinal) && request != "GetData /" + && !request.EndsWith("/IAmAlive", StringComparison.Ordinal) ? release.Task : Task.CompletedTask; + var read = harness.Read(cancellationToken: cancellation.Token); + var owner = Assert.Single(harness.Sessions); + try + { + Assert.Equal( + new[] { "Sync /", "GetChildren /", "GetData " + ZooKeeperNativeFake.RowPath(harness.Entries[0].SiloAddress) }, + harness.Calls); + if (cancel) + { + cancellation.Cancel(); + await Assert.ThrowsAnyAsync(() => read); + } + Assert.False(owner.Completion.IsCompleted); + Assert.Equal(0, harness.CloseCount); + } + finally + { + release.TrySetResult(); + await Record.ExceptionAsync(() => owner.Completion); + } + if (cancel) + { + var failure = await Assert.ThrowsAnyAsync(() => owner.Completion); + Assert.Equal(cancellation.Token, failure.CancellationToken); + Assert.Equal(3, harness.Calls.Count); + } + else + { + harness.AssertSnapshot(await read, harness.Entries, rowCount); + Assert.Equal(3 + rowCount * 2, harness.Calls.Count); + } + harness.AssertOneOwner(readOnly: true); + } + + [Fact] + public async Task Read_MissingRow_StopsBeforeLaterNativeFailure() + { + var harness = await Harness.CreateAsync(); + var missing = new KeeperException.NoNodeException(ZooKeeperNativeFake.RowPath(harness.Entries[0].SiloAddress)); + var laterFailure = new KeeperException.NoAuthException(); + var first = "GetData " + ZooKeeperNativeFake.RowPath(harness.Entries[0].SiloAddress); + var second = "GetData " + ZooKeeperNativeFake.RowPath(harness.Entries[1].SiloAddress); + harness.BeforeRequest = request => + { + if (request == first) + throw missing; + if (request == second) + throw laterFailure; + return Task.CompletedTask; + }; + + Assert.Same(missing, await Record.ExceptionAsync(() => harness.Read(TestContext.Current.CancellationToken))); + Assert.DoesNotContain(second, harness.Calls); + Assert.Contains("GetData /", harness.Calls); + harness.AssertOneOwner(readOnly: true); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task GetGateways_ConnectionLoss_UsesRetriedSnapshot(bool exhaust) + { + var harness = await Harness.CreateAsync(3); + harness.Entries[1].ProxyPort = 0; + harness.Entries[2].Status = SiloStatus.Dead; + foreach (var entry in harness.Entries) + { + var path = ZooKeeperNativeFake.RowPath(entry.SiloAddress); + harness.Fake.Nodes[path] = harness.Fake.Nodes[path] with { Data = ZooKeeperBasedMembershipTable.Serialize(entry) }; + } + var failure = new KeeperException.ConnectionLossException(); + var failures = 0; + harness.BeforeRequest = request => + { + if (request == "Sync /" && (exhaust || failures++ == 0)) + throw failure; + return Task.CompletedTask; + }; + var provider = new ZooKeeperGatewayListProvider( + NullLogger.Instance, + Options.Create(new ZooKeeperGatewayListProviderOptions { ConnectionString = "unused.invalid" }), + Options.Create(new GatewayOptions()), + Options.Create(new ClusterOptions { ClusterId = "test" }), + () => harness.CreateSession(true), + harness.Pipeline); + var read = provider.GetGateways(); + var completion = Record.ExceptionAsync(() => read); + foreach (var delay in exhaust ? new[] { 250, 500, 1000, 2000 } : [250]) + await harness.Clock.AdvanceNextAsync(TimeSpan.FromMilliseconds(delay)); + if (exhaust) + { + Assert.Same(failure, await completion); + } + else + { + Assert.Null(await completion); + var entry = harness.Entries[0]; + Assert.Equal(SiloAddress.New(entry.SiloAddress.Endpoint.Address, entry.ProxyPort, entry.SiloAddress.Generation).ToGatewayUri(), + Assert.Single(await read)); + } + harness.AssertOneOwner(readOnly: true); + } + + [Theory] + [InlineData(false, false)] + [InlineData(false, true)] + [InlineData(true, false)] + [InlineData(true, true)] + public async Task ConditionalWrite_CommitOrCloseLoss_PropagatesWithoutReplay(bool update, bool closeFailure) + { + var harness = await Harness.CreateAsync(1); + var entry = update ? harness.Entries[0] : Harness.Entry(1); + entry.Status = SiloStatus.Dead; + var failure = new KeeperException.ConnectionLossException(); + if (closeFailure) + harness.Close = () => Task.FromException(failure); + else + harness.AfterMulti = () => Task.FromException(failure); + + var actual = await Record.ExceptionAsync(() => update + ? harness.Table.UpdateRowAsync(entry, "0", new TableVersion(2, "1"), TestContext.Current.CancellationToken) + : harness.Table.InsertRowAsync(entry, new TableVersion(2, "1"), TestContext.Current.CancellationToken)); + + Assert.Same(failure, actual); + Assert.Equal("Multi", Assert.Single(harness.Calls)); + Assert.Single(harness.Fake.Transactions); + Assert.Equal(2, harness.Fake.Nodes["/"].Version); + var row = harness.Fake.Nodes[ZooKeeperNativeFake.RowPath(entry.SiloAddress)]; + Assert.Equal(update ? 1 : 0, row.Version); + Assert.Equal(ZooKeeperBasedMembershipTable.Serialize(entry), row.Data); + Assert.Equal(update ? 3 : 5, harness.Fake.Nodes.Count); + Assert.Empty(harness.Logger.Warnings); + harness.AssertOneOwner(readOnly: false); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task ConditionalWrite_KnownConflict_PreservesFalseResult(bool update) + { + var harness = await Harness.CreateAsync(1); + var before = harness.Fake.Nodes.ToDictionary(pair => pair.Key, pair => pair.Value); + var result = update + ? await harness.Table.UpdateRowAsync(harness.Entries[0], "99", new TableVersion(2, "1"), TestContext.Current.CancellationToken) + : await harness.Table.InsertRowAsync(harness.Entries[0], new TableVersion(2, "1"), TestContext.Current.CancellationToken); + + Assert.False(result); + Assert.Equal(before, harness.Fake.Nodes); + Assert.Equal("Multi", Assert.Single(harness.Calls)); + Assert.Empty(harness.Logger.Warnings); + harness.AssertOneOwner(readOnly: false); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task ConditionalWrite_CallbackAndCloseFailures_PreserveBothExceptions(bool update) + { + var harness = await Harness.CreateAsync(1); + var entry = update ? harness.Entries[0] : Harness.Entry(1); + entry.Status = SiloStatus.Dead; + var primary = new KeeperException.ConnectionLossException(); + var secondary = new KeeperException.ConnectionLossException(); + harness.AfterMulti = () => Task.FromException(primary); + harness.Close = () => Task.FromException(secondary); + + var actual = await Assert.ThrowsAsync(() => update + ? harness.Table.UpdateRowAsync(entry, "0", new TableVersion(2, "1"), TestContext.Current.CancellationToken) + : harness.Table.InsertRowAsync(entry, new TableVersion(2, "1"), TestContext.Current.CancellationToken)); + + Assert.Equal(new Exception[] { primary, secondary }, actual.InnerExceptions); + Assert.Same(actual, await Record.ExceptionAsync(() => Assert.Single(harness.Sessions).Completion)); + Assert.Equal("Multi", Assert.Single(harness.Calls)); + Assert.Single(harness.Fake.Transactions); + Assert.Equal(2, harness.Fake.Nodes["/"].Version); + var row = harness.Fake.Nodes[ZooKeeperNativeFake.RowPath(entry.SiloAddress)]; + Assert.Equal(update ? 1 : 0, row.Version); + Assert.Equal(ZooKeeperBasedMembershipTable.Serialize(entry), row.Data); + Assert.Empty(harness.Logger.Warnings); + harness.AssertOneOwner(readOnly: false); + } + + [Fact] + public async Task ReadDecorator_PreservesMutationDelegates() + { + var harness = await Harness.CreateAsync(); + var native = harness.Fake.Operations; + var wrapped = ZooKeeperReadRetryPolicy.Wrap(native, harness.Pipeline, CancellationToken.None); + + Assert.Same(native.Multi, wrapped.Multi); + Assert.Same(native.SetData, wrapped.SetData); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task NativeFixture_TeardownTimeout_JoinsActualCompletion(bool fail) + { + var nativeCompletion = Gate(); + var timeout = new TimeoutException("fixture teardown wait expired"); + timeout.Data["ClusteringTestKit.CleanupCompletion"] = nativeCompletion.Task; + var nativeFailure = new KeeperException.ConnectionLossException(); + var drain = ZooKeeperReadResilienceTests.DrainFixtureAsync(() => ValueTask.FromException(timeout)); + + try + { + Assert.False(drain.IsCompleted); + } + finally + { + if (fail) + nativeCompletion.SetException(nativeFailure); + else + nativeCompletion.SetResult(); + } + + Assert.Same(fail ? (Exception)nativeFailure : timeout, await Record.ExceptionAsync(() => drain)); + Assert.True(nativeCompletion.Task.IsCompleted); + } + + private static TaskCompletionSource Gate() => new(TaskCreationOptions.RunContinuationsAsynchronously); + + private static ZooKeeperWatcher CreateConnectionMonitor( + TimeProvider? timeProvider = null, + TimeSpan? reconnectTimeout = null) => + new( + NullLogger.Instance, + timeProvider ?? TimeProvider.System, + reconnectTimeout ?? TimeSpan.FromMinutes(1)); + + private static Task Signal(ZooKeeperWatcher watcher, Watcher.Event.KeeperState state) + { + watcher.ProcessConnectionState(state); + return Task.CompletedTask; + } + + [Fact] + public async Task NativeFixture_PrimaryFailure_IsCapturedBeforeTeardownAndRethrownUnchanged() + { + var harness = await Harness.CreateAsync(); + var primary = new KeeperException.NoAuthException(); + harness.BeforeRequest = _ => Task.FromException(primary); + var events = new List(); + string captured = null!; + async Task PrimaryReadScenario() + { + await harness.Read(TestContext.Current.CancellationToken); + } + async Task RunWithTeardown() + { + try + { + await ZooKeeperReadResilienceTests.CapturePrimaryScenarioFailureAsync(PrimaryReadScenario, text => + { + events.Add("capture"); + Assert.Equal(primary.ToString(), text); + captured = text; + }); + } + finally + { + events.Add("teardown"); + Assert.Same(primary, await Record.ExceptionAsync(() => Assert.Single(harness.Sessions).Completion)); + } + } + + Assert.Same(primary, await Record.ExceptionAsync(RunWithTeardown)); + Assert.Equal(new[] { "capture", "teardown" }, events); + Assert.Contains(nameof(KeeperException.NoAuthException), captured, StringComparison.Ordinal); + Assert.Contains(nameof(PrimaryReadScenario), captured, StringComparison.Ordinal); + harness.AssertOneOwner(readOnly: true); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task NativeFixture_CanceledCleanup_PreservesAdmissionAndOwnsClose(bool fail) + { + var harness = await Harness.CreateAsync(); + foreach (var entry in harness.Entries) + { + entry.Status = SiloStatus.Dead; + var path = ZooKeeperNativeFake.RowPath(entry.SiloAddress); + harness.Fake.Nodes[path] = harness.Fake.Nodes[path] with { Data = ZooKeeperBasedMembershipTable.Serialize(entry) }; + } + var first = Gate(); + var firstSettled = Gate(); + var closeStarted = Gate(); + var releaseClose = Gate(); + var failure = new KeeperException.NoAuthException(); + var firstPath = ZooKeeperNativeFake.RowPath(harness.Entries[0].SiloAddress); + var secondPath = ZooKeeperNativeFake.RowPath(harness.Entries[1].SiloAddress); + harness.Fake.BeforeRead = async path => + { + if (path == firstPath) + { + try + { + await first.Task; + if (fail) + throw failure; + } + finally + { + firstSettled.TrySetResult(); + } + } + }; + using var cancellation = new CancellationTokenSource(); + async Task NativeOperation() + { + try + { + return await ZooKeeperBasedMembershipTable.CleanupCoreAsync( + harness.Fake.Operations, DateTimeOffset.UnixEpoch.AddDays(3), cancellation.Token); + } + finally + { + closeStarted.SetResult(); + await releaseClose.Task; + } + } + var native = NativeOperation(); + Task? observed = null; + var caller = ZooKeeperBasedMembershipTable.AwaitOperationAsync(native, cancellation.Token, operation => observed = operation); + try + { + Assert.Same(native, observed); + Assert.Equal(new[] { "children /", "read " + firstPath }, harness.Fake.Calls); + cancellation.Cancel(); + var canceled = await Assert.ThrowsAnyAsync(() => caller); + Assert.Equal(cancellation.Token, canceled.CancellationToken); + Assert.False(observed!.IsCompleted); + first.SetResult(); + await firstSettled.Task.WaitAsync(TestContext.Current.CancellationToken); + await closeStarted.Task.WaitAsync(TestContext.Current.CancellationToken); + Assert.False(observed.IsCompleted); + Assert.DoesNotContain("read " + secondPath, harness.Fake.Calls); + } + finally + { + first.TrySetResult(); + releaseClose.TrySetResult(); + await Record.ExceptionAsync(() => native); + } + var actual = await Record.ExceptionAsync(() => observed!); + if (fail) + Assert.Same(failure, actual); + else + Assert.Equal(cancellation.Token, Assert.IsAssignableFrom(actual).CancellationToken); + Assert.Empty(harness.Fake.Transactions); + Assert.Equal(new[] { "children /", "read " + firstPath }, harness.Fake.Calls); + } + + private sealed class Harness + { + internal ZooKeeperNativeFake Fake { get; } = new(); + internal RetryClock Clock { get; } = new(); + internal RecordingLogger Logger { get; } = new(); + internal ConcurrentQueue Calls { get; } = new(); + internal List Sessions { get; } = []; + internal List ReadOnly { get; } = []; + internal Action? OnSessionCreated { get; set; } + internal IZooKeeperConnectionMonitor? ConnectionMonitor { get; set; } + internal Func? BeforeRequest { get; set; } + internal Action? AfterRequest { get; set; } + internal Func Close { get; set; } = () => Task.CompletedTask; + internal Func AfterMulti { get; set; } = () => Task.CompletedTask; + internal int CloseCount; + internal MembershipEntry[] Entries { get; private set; } = []; + internal ResiliencePipeline Pipeline { get; } + internal ZooKeeperBasedMembershipTable Table { get; } + + private Harness() + { + Pipeline = ZooKeeperReadRetryPolicy.CreatePipeline(Logger, Clock); + Table = new ZooKeeperBasedMembershipTable(NullLogger.Instance, + Options.Create(new ZooKeeperClusteringSiloOptions { ConnectionString = "unused.invalid" }), + Options.Create(new ClusterOptions { ClusterId = "test" }), CreateSession, Pipeline); + } + + internal static async Task CreateAsync(int rowCount = 2) + { + var result = new Harness { Entries = Enumerable.Range(0, rowCount).Select(Entry).ToArray() }; + for (var index = 0; index < rowCount; index++) + { + Assert.True(await ZooKeeperBasedMembershipTable.InsertRowCoreAsync(result.Fake.Operations, + result.Entries[index], new TableVersion(index + 1, index.ToString(CultureInfo.InvariantCulture)), + TestContext.Current.CancellationToken)); + } + result.Fake.Calls.Clear(); + result.Fake.Transactions.Clear(); + return result; + } + + internal static MembershipEntry Entry(int index) => new() + { + SiloAddress = SiloAddress.New(new IPEndPoint(IPAddress.Loopback, 11111), 12345 + index), + HostName = "host-a", + SiloName = "silo-a", + Status = SiloStatus.Active, + ProxyPort = 30000 + index, + StartTime = DateTime.UnixEpoch, + IAmAliveTime = DateTime.UnixEpoch.AddDays(1) + }; + + internal ZooKeeperSession CreateSession(bool readOnly) + { + ReadOnly.Add(readOnly); + var native = Fake.Operations; + var operations = new ZooKeeperBasedMembershipTable.NativeOperations( + path => Request("GetData " + path, () => native.GetData(path)), + path => Request("GetChildren " + path, () => native.GetChildren(path)), + path => Request("Sync " + path, async () => { await native.Sync(path); return true; }), + async ops => + { + Calls.Enqueue("Multi"); + await native.Multi(ops); + await AfterMulti(); + }, + native.SetData); + var session = new ZooKeeperSession( + operations, + () => + { + Interlocked.Increment(ref CloseCount); + return Close(); + }, + ConnectionMonitor); + Sessions.Add(session); + OnSessionCreated?.Invoke(session); + return session; + } + + private async Task Request(string request, Func> action) + { + Calls.Enqueue(request); + if (BeforeRequest is { } before) + await before(request); + Task native; + lock (Fake.Calls) + native = action(); + var result = await native; + AfterRequest?.Invoke(request); + return result; + } + + internal Task Read(CancellationToken cancellationToken) => Read(false, cancellationToken); + + internal Task Read(bool point, CancellationToken cancellationToken) => + point ? Table.ReadRowAsync(Entries[0].SiloAddress, cancellationToken) : Table.ReadAllAsync(cancellationToken); + + internal IEnumerable ExpectedReadCalls(bool point = false) + { + yield return "Sync /"; + yield return point ? "GetData /" : "GetChildren /"; + foreach (var entry in point ? Entries.Take(1) : Entries) + { + yield return "GetData " + ZooKeeperNativeFake.RowPath(entry.SiloAddress); + yield return "GetData " + ZooKeeperNativeFake.HeartbeatPath(entry.SiloAddress); + } + yield return "GetData /"; + } + + internal Dictionary CountCalls() => + Calls.GroupBy(value => value).ToDictionary(group => group.Key, group => group.Count()); + + internal void AssertSnapshot(MembershipTableData snapshot, IEnumerable expected, int version) + { + Assert.Equal(version, snapshot.Version.Version); + Assert.Equal(version.ToString(CultureInfo.InvariantCulture), snapshot.Version.VersionEtag); + var entries = expected.ToDictionary(entry => entry.SiloAddress); + var actual = snapshot.Members.ToDictionary(row => row.Item1.SiloAddress); + Assert.Equal(entries.Count, actual.Count); + foreach (var (address, entry) in entries) + { + var row = actual[address]; + Assert.Equal(Fake.Nodes[ZooKeeperNativeFake.RowPath(address)].Version.ToString(CultureInfo.InvariantCulture), row.Item2); + Assert.Equal(ZooKeeperBasedMembershipTable.Serialize(entry), ZooKeeperBasedMembershipTable.Serialize(row.Item1)); + } + } + + internal void AssertOneOwner(bool readOnly) + { + Assert.Single(Sessions); + Assert.Equal(readOnly, Assert.Single(ReadOnly)); + Assert.Equal(1, CloseCount); + Assert.True(Sessions[0].Completion.IsCompleted); + } + } + + private sealed class RetryClock : TimeProvider + { + private readonly FakeTimeProvider _time = new(); + private readonly Channel _timers = Channel.CreateUnbounded(); + public override DateTimeOffset GetUtcNow() => _time.GetUtcNow(); + public override long GetTimestamp() => _time.GetTimestamp(); + public override long TimestampFrequency => _time.TimestampFrequency; + + public override ITimer CreateTimer(TimerCallback callback, object? state, TimeSpan dueTime, TimeSpan period) + { + var timer = _time.CreateTimer(callback, state, dueTime, period); + Assert.True(_timers.Writer.TryWrite(dueTime)); + return timer; + } + + internal ValueTask NextTimerAsync() => _timers.Reader.ReadAsync(TestContext.Current.CancellationToken); + internal void Advance(TimeSpan duration) => _time.Advance(duration); + internal async Task AdvanceNextAsync(TimeSpan expected) + { + Assert.Equal(expected, await NextTimerAsync()); + Advance(expected); + } + } + + private sealed class RecordingLogger : ILogger + { + internal sealed record Warning(Exception? Exception, IReadOnlyDictionary Values); + internal ConcurrentQueue Warnings { get; } = new(); + public IDisposable? BeginScope(TState state) where TState : notnull => null; + public bool IsEnabled(LogLevel logLevel) => true; + public void Log(LogLevel logLevel, EventId eventId, TState state, Exception? exception, Func formatter) + { + Assert.Equal(LogLevel.Warning, logLevel); + var values = Assert.IsAssignableFrom>>(state); + Warnings.Enqueue(new(exception, values.ToDictionary(pair => pair.Key, pair => pair.Value))); + } + } +}