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)));
+ }
+ }
+}