diff --git a/docs/site/src/content/docs/host/configuration-guide/local-development-configuration.md b/docs/site/src/content/docs/host/configuration-guide/local-development-configuration.md index a35c37f5148..41b4d1bd4a2 100644 --- a/docs/site/src/content/docs/host/configuration-guide/local-development-configuration.md +++ b/docs/site/src/content/docs/host/configuration-guide/local-development-configuration.md @@ -33,6 +33,8 @@ The client and silo must use matching gateway ports, . For overlapping reads, the comparison uses the version known when each read began. + Aspire is usually the easiest option because it allocates endpoints, starts dependencies, injects configuration, and displays logs for every replica. > [!IMPORTANT] diff --git a/docs/site/src/content/docs/implementation/cluster-management.md b/docs/site/src/content/docs/implementation/cluster-management.md index 1e8b1a7190b..3656d523191 100644 --- a/docs/site/src/content/docs/implementation/cluster-management.md +++ b/docs/site/src/content/docs/implementation/cluster-management.md @@ -11,9 +11,15 @@ Cluster membership answers one question for the rest of the runtime: which silo ## Identity, status, and views -A silo identity includes its advertised endpoint and a generation value, so a restarted process at the same endpoint is a new identity. Its progresses through `Created`, `Joining`, `Active`, and a terminating status (`ShuttingDown`, `Stopping`, or `Dead`). +A silo identity includes its advertised endpoint and a generation value, so a restarted process at the same endpoint is a new identity. Its progresses in one direction through `Created`, `Joining`, `Active`, and terminating states (`ShuttingDown`, `Stopping`, or `Dead`). -Each successful versioned membership-table mutation, such as inserting a row or changing a status, advances the table version. Periodic writes leave the version unchanged. `MembershipTableManager` publishes immutable snapshots through `ClusterMembershipService`; consumers ignore older versions. Directory ownership, gateway discovery, and failure recovery therefore observe a monotonically ordered sequence of views even if notifications arrive out of order. +The membership system provides a **canonical membership view**: a versioned view of silo identities and their states. Each successful versioned membership-table mutation advances the version by exactly one. A version uniquely identifies its canonical membership view throughout the cluster, so two hosts observing version N have the same canonical membership view. A host can advance directly from N to N + 2 when its next refresh observes two completed mutations. + +`MembershipTableManager` publishes snapshots through `ClusterMembershipService` with monotonically advancing canonical membership views. Directory ownership, gateway discovery, and failure recovery rely on this guarantee. + +Per-silo `IAmAliveTime` is tracked independently of the canonical membership view. Periodic writes leave the view version unchanged. Snapshot updates retain the maximum observed timestamp for each silo, so local liveness timestamps advance monotonically as table reads and peer snapshots arrive. At the same version, merging retains the accepted versioned fields. + +Snapshots can prune previously `Dead` rows at the same version while retaining every non-Dead row. Pruning preserves the versioned fields and maximum `IAmAliveTime` of each retained entry. ```mermaid flowchart LR diff --git a/src/Orleans.Core/Runtime/MembershipTableSnapshot.cs b/src/Orleans.Core/Runtime/MembershipTableSnapshot.cs index 8b14f96fc29..d059588caad 100644 --- a/src/Orleans.Core/Runtime/MembershipTableSnapshot.cs +++ b/src/Orleans.Core/Runtime/MembershipTableSnapshot.cs @@ -7,8 +7,12 @@ namespace Orleans.Runtime { /// - /// Represents an immutable snapshot of cluster membership state. + /// Represents an immutable snapshot of a canonical membership view and per-silo liveness timestamps. /// + /// + /// A version identifies the same canonical membership view throughout the cluster. + /// Updates retain versioned fields at the same version and the maximum observed IAmAliveTime for each silo. + /// [GenerateSerializer, Immutable] internal sealed class MembershipTableSnapshot : ISpanFormattable { @@ -63,33 +67,35 @@ public static MembershipTableSnapshot Update(MembershipTableSnapshot previousSna return Update(previousSnapshot, updated.Version, updated.Entries.Values); } - private static MembershipTableSnapshot Update(MembershipTableSnapshot previousSnapshot, MembershipVersion version, IEnumerable updatedEntries) + private static MembershipTableSnapshot Update( + MembershipTableSnapshot previousSnapshot, + MembershipVersion version, + IEnumerable updatedEntries) { - ArgumentNullException.ThrowIfNull(previousSnapshot); - ArgumentNullException.ThrowIfNull(updatedEntries); - var entries = ImmutableDictionary.CreateBuilder(); foreach (var item in updatedEntries) { var entry = item; - entry = PreserveIAmAliveTime(previousSnapshot, entry); - entries.Add(entry.SiloAddress, entry); - } + if (previousSnapshot.Entries.TryGetValue(entry.SiloAddress, out var previousEntry)) + { + var iAmAliveTime = entry.IAmAliveTime > previousEntry.IAmAliveTime + ? entry.IAmAliveTime + : previousEntry.IAmAliveTime; + if (version == previousSnapshot.Version) + { + entry = previousEntry; + } - return new MembershipTableSnapshot(version, entries.ToImmutable()); - } + if (entry.IAmAliveTime < iAmAliveTime) + { + entry = entry.WithIAmAliveTime(iAmAliveTime); + } + } - private static MembershipEntry PreserveIAmAliveTime(MembershipTableSnapshot previousSnapshot, MembershipEntry entry) - { - // Retain the maximum IAmAliveTime, since IAmAliveTime updates do not increase membership version - // and therefore can be clobbered by torn reads. - if (previousSnapshot.Entries.TryGetValue(entry.SiloAddress, out var previousEntry) - && previousEntry.IAmAliveTime > entry.IAmAliveTime) - { - entry = entry.WithIAmAliveTime(previousEntry.IAmAliveTime); + entries.Add(entry.SiloAddress, entry); } - return entry; + return new MembershipTableSnapshot(version, entries.ToImmutable()); } /// @@ -150,6 +156,10 @@ public SiloStatus GetSiloStatus(SiloAddress silo) /// /// Determines whether this snapshot is a successor to another snapshot. /// + /// + /// At the same canonical membership version, progress consists of newer liveness timestamps + /// or pruning previously Dead rows. Non-Dead rows are retained. + /// /// The snapshot to compare against. /// if this snapshot is a successor to ; otherwise, . public bool IsSuccessorTo(MembershipTableSnapshot other) @@ -166,29 +176,36 @@ public bool IsSuccessorTo(MembershipTableSnapshot other) if (Entries.Count > other.Entries.Count) { - // Something is amiss. + // Adding a row requires a new canonical membership version. return false; } - foreach (var entry in Entries) + var heartbeatAdvanced = false; + foreach (var (silo, entry) in Entries) { - if (!other.Entries.TryGetValue(entry.Key, out var otherEntry)) + if (!other.Entries.TryGetValue(silo, out var otherEntry)) { - // Something is amiss. + // Membership changes require a table-version advance. return false; } + + heartbeatAdvanced |= entry.IAmAliveTime > otherEntry.IAmAliveTime; } - // This is a successor if any silo has a later EffectiveIAmAliveTime. - foreach (var entry in Entries) + if (Entries.Count == other.Entries.Count) { - if (entry.Value.EffectiveIAmAliveTime > other.Entries[entry.Key].EffectiveIAmAliveTime) + return heartbeatAdvanced; + } + + foreach (var (silo, previousEntry) in other.Entries) + { + if (previousEntry.Status != SiloStatus.Dead && !Entries.ContainsKey(silo)) { - return true; + return false; } } - return false; + return true; } public override string ToString() diff --git a/src/Orleans.Runtime/MembershipService/MembershipTableManager.cs b/src/Orleans.Runtime/MembershipService/MembershipTableManager.cs index df986d75f3c..9d08d66ae62 100644 --- a/src/Orleans.Runtime/MembershipService/MembershipTableManager.cs +++ b/src/Orleans.Runtime/MembershipService/MembershipTableManager.cs @@ -186,6 +186,7 @@ public async Task RefreshFromSnapshot(MembershipTableSnapshot snapshot, Cancella LogInformationReceivedClusterMembershipSnapshot(this.log, snapshot); + cancellationToken.ThrowIfCancellationRequested(); this.TryProcessMembershipUpdate(MembershipTableSnapshot.Update, snapshot, nameof(RefreshFromSnapshot)); } @@ -562,6 +563,11 @@ private MembershipTableSnapshot ProcessMembershipUpdate( MembershipTableSnapshot previous, MembershipTableSnapshot updated) { + if (updated.Version < previous.Version) + { + return previous; + } + if (!previous.Entries.TryGetValue(this.myAddress, out var previousLocalSiloEntry) || previousLocalSiloEntry.Status == SiloStatus.Created) { diff --git a/src/Orleans.Runtime/MembershipService/SystemTargetBasedMembershipTable.cs b/src/Orleans.Runtime/MembershipService/SystemTargetBasedMembershipTable.cs index 91438bdef61..324bfb2ffb8 100644 --- a/src/Orleans.Runtime/MembershipService/SystemTargetBasedMembershipTable.cs +++ b/src/Orleans.Runtime/MembershipService/SystemTargetBasedMembershipTable.cs @@ -15,13 +15,21 @@ internal partial class SystemTargetBasedMembershipTable : IMembershipTable { private readonly IServiceProvider serviceProvider; private readonly ILogger logger; + private readonly IFatalErrorHandler fatalErrorHandler; + private IMembershipManager? membershipManager; + private int fatalErrorReported; private IMembershipTableSystemTarget grain = null!; - public SystemTargetBasedMembershipTable(IServiceProvider serviceProvider, ILogger logger) + public SystemTargetBasedMembershipTable( + IServiceProvider serviceProvider, + ILogger logger, + IFatalErrorHandler fatalErrorHandler) { this.serviceProvider = serviceProvider; this.logger = logger; + this.fatalErrorHandler = fatalErrorHandler; } + [Obsolete("Use InitializeMembershipTableAsync instead.")] public Task InitializeMembershipTable(bool tryInitTableVersion) => InitializeMembershipTableAsync(tryInitTableVersion, CancellationToken.None); @@ -103,12 +111,52 @@ private async Task WaitForTableGrainToInit(IMembershipTableSystemTarget membersh [Obsolete("Use ReadRowAsync instead.")] public Task ReadRow(SiloAddress key) => ReadRowAsync(key, CancellationToken.None); - public Task ReadRowAsync(SiloAddress key, CancellationToken cancellationToken = default) => this.grain.ReadRowAsync(key, cancellationToken); + public async Task ReadRowAsync(SiloAddress key, CancellationToken cancellationToken = default) + { + cancellationToken.ThrowIfCancellationRequested(); + var observedVersion = GetCurrentVersion(); + var table = await this.grain.ReadRowAsync(key, cancellationToken); + cancellationToken.ThrowIfCancellationRequested(); + ValidateVersion(observedVersion, table); + return table; + } [Obsolete("Use ReadAllAsync instead.")] public Task ReadAll() => ReadAllAsync(CancellationToken.None); - public Task ReadAllAsync(CancellationToken cancellationToken = default) => this.grain.ReadAllAsync(cancellationToken); + public async Task ReadAllAsync(CancellationToken cancellationToken = default) + { + cancellationToken.ThrowIfCancellationRequested(); + var observedVersion = GetCurrentVersion(); + var table = await this.grain.ReadAllAsync(cancellationToken); + cancellationToken.ThrowIfCancellationRequested(); + ValidateVersion(observedVersion, table); + return table; + } + + private MembershipVersion GetCurrentVersion() + { + // Resolve after owner construction, and include versions learned through committed writes and gossip. + var manager = this.membershipManager ??= this.serviceProvider.GetRequiredService(); + return manager.CurrentSnapshot.Version; + } + + private void ValidateVersion(MembershipVersion observedVersion, MembershipTableData table) + { + // Compare with the version known before this read began so overlapping reads can finish out of order. + if (table.Version.Version < observedVersion.Value) + { + var reason = $"The development membership table version decreased from {observedVersion} to {table.Version.Version}. " + + "The cluster's membership state has been lost and this silo must terminate."; + var exception = new OrleansException(reason); + if (Interlocked.Exchange(ref this.fatalErrorReported, 1) == 0) + { + this.fatalErrorHandler.OnFatalException(this, reason, exception); + } + + throw exception; + } + } [Obsolete("Use InsertRowAsync instead.")] public Task InsertRow(MembershipEntry entry, TableVersion tableVersion) => InsertRowAsync(entry, tableVersion, CancellationToken.None); diff --git a/test/Orleans.Core.Tests/Membership/MembershipTableManagerTests.cs b/test/Orleans.Core.Tests/Membership/MembershipTableManagerTests.cs index eb9f417c3bd..e4f884e0f92 100644 --- a/test/Orleans.Core.Tests/Membership/MembershipTableManagerTests.cs +++ b/test/Orleans.Core.Tests/Membership/MembershipTableManagerTests.cs @@ -1,5 +1,6 @@ using System.Collections.Concurrent; using System.Collections.Immutable; +using System.Threading.Channels; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using Microsoft.Extensions.Time.Testing; @@ -32,6 +33,7 @@ public class MembershipTableManagerTests private readonly IFatalErrorHandler fatalErrorHandler; private readonly IMembershipGossiper membershipGossiper; private readonly SiloLifecycleSubject lifecycle; + private MembershipTableManager? membershipTableManager; public MembershipTableManagerTests(ITestOutputHelper output) { @@ -1331,6 +1333,48 @@ await Assert.ThrowsAnyAsync( Assert.Equal(new MembershipVersion(1), manager.MembershipTableSnapshot.Version); } + [Fact] + public async Task MembershipOwnerRechecksCancellationAfterSharedRefreshCompleted() + { + var read = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var table = Substitute.For(); + table.ReadAllAsync(Arg.Any()).Returns(read.Task); + using var manager = CreateMembershipTableManager(table); + var refresh = manager.Refresh(cancellationToken: TestContext.Current.CancellationToken); + var context = new MembershipContinuationContext(); + using var cancellation = CancellationTokenSource.CreateLinkedTokenSource(TestContext.Current.CancellationToken); + var peer = Silo("127.0.0.1:200@100"); + var incoming = Snapshot( + new MembershipVersion(2), Entry(peer, SiloStatus.Active, DateTimeOffset.UnixEpoch)); + Task gossip; + var previousContext = SynchronizationContext.Current; + SynchronizationContext.SetSynchronizationContext(context); + try + { + gossip = ((IMembershipManager)manager).ProcessGossipSnapshot(incoming, cancellation.Token); + } + finally + { + SynchronizationContext.SetSynchronizationContext(previousContext); + } + + Assert.False(gossip.IsCompleted); + read.SetResult(new MembershipTableData(new TableVersion(1, "1"))); + var continuation = await context.TakeContinuation(TestContext.Current.CancellationToken); + await refresh.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken); + var committed = manager.MembershipTableSnapshot; + Assert.Equal(new MembershipVersion(1), committed.Version); + + cancellation.Cancel(); + continuation.Callback(continuation.State); + + var exception = await Assert.ThrowsAnyAsync( + () => gossip.WaitAsync(TimeSpan.FromSeconds(5), TestContext.Current.CancellationToken)); + Assert.Equal(cancellation.Token, exception.CancellationToken); + Assert.Same(committed, manager.MembershipTableSnapshot); + Assert.DoesNotContain(peer, manager.MembershipTableSnapshot.Entries.Keys); + } + [Theory] [InlineData(false, 3000)] [InlineData(true, 500)] @@ -1453,6 +1497,8 @@ private SystemTargetBasedMembershipTable CreateSystemTargetBasedMembershipTable( var membershipTarget = Substitute.For(); membershipTarget.ReadAllAsync(Arg.Any()) .Returns(call => backingTable.ReadAllAsync(call.ArgAt(0))); + membershipTarget.ReadRowAsync(Arg.Any(), Arg.Any()) + .Returns(call => backingTable.ReadRowAsync(call.ArgAt(0), call.ArgAt(1))); membershipTarget.InsertRowAsync(Arg.Any(), Arg.Any(), Arg.Any()) .Returns(call => backingTable.InsertRowAsync(call.ArgAt(0), call.ArgAt(1), call.ArgAt(2))); membershipTarget.UpdateRowAsync(Arg.Any(), Arg.Any(), Arg.Any(), Arg.Any()) @@ -1465,7 +1511,216 @@ private SystemTargetBasedMembershipTable CreateSystemTargetBasedMembershipTable( Options.Create(new DevelopmentClusterMembershipOptions { PrimarySiloEndpoint = primarySilo.Endpoint })); services.GetService(typeof(ILocalSiloDetails)).Returns(this.localSiloDetails); services.GetService(typeof(IInternalGrainFactory)).Returns(grainFactory); - return new SystemTargetBasedMembershipTable(services, this.loggerFactory.CreateLogger()); + services.GetService(typeof(IMembershipManager)).Returns(_ => this.membershipTableManager); + return new SystemTargetBasedMembershipTable( + services, this.loggerFactory.CreateLogger(), this.fatalErrorHandler); + } + + [Theory] + [InlineData(100, false)] + [InlineData(101, false)] + [InlineData(100, true)] + [InlineData(101, true)] + public async Task DevelopmentMembershipReadChecksVersionLearnedFromGossip(int tableVersion, bool readRow) + { + var cancellationToken = TestContext.Current.CancellationToken; + var local = Entry(this.localSilo, SiloStatus.Active, DateTimeOffset.UnixEpoch); + var backing = Substitute.For(); + backing.ReadAllAsync(Arg.Any()).Returns( + new MembershipTableData(Tuple.Create(local, "local"), new TableVersion(100, "100"))); + var provider = CreateSystemTargetBasedMembershipTable(backing); + await provider.InitializeMembershipTableAsync(true, cancellationToken); + using var manager = CreateMembershipTableManager(provider); + await manager.Refresh(cancellationToken: cancellationToken); + backing.ClearReceivedCalls(); + var peer = Entry(Silo("127.0.0.1:200@100"), SiloStatus.Active, DateTimeOffset.UnixEpoch); + var incoming = Snapshot(new MembershipVersion(101), local, peer); + + await manager.RefreshFromSnapshot(incoming, cancellationToken); + + var accepted = manager.MembershipTableSnapshot; + Assert.Equal(new MembershipVersion(101), accepted.Version); + Assert.Contains(peer.SiloAddress, accepted.Entries.Keys); + Assert.Empty(backing.ReceivedCalls()); + this.fatalErrorHandler.DidNotReceiveWithAnyArgs().OnFatalException(default, default, default); + var version = new TableVersion(tableVersion, tableVersion.ToString()); + var entries = new List> { Tuple.Create(local, "local") }; + if (tableVersion == 101) + { + entries.Add(Tuple.Create(peer, "peer")); + } + + backing.ReadAllAsync(Arg.Any()).Returns(new MembershipTableData(entries, version)); + backing.ReadRowAsync(this.localSilo, Arg.Any()) + .Returns(new MembershipTableData(Tuple.Create(local, "local"), version)); + var read = readRow + ? provider.ReadRowAsync(this.localSilo, cancellationToken) + : provider.ReadAllAsync(cancellationToken); + if (tableVersion < accepted.Version.Value) + { + var error = await Assert.ThrowsAsync(() => read); + Assert.Contains("version decreased from 101 to 100", error.Message, StringComparison.Ordinal); + this.fatalErrorHandler.Received(1).OnFatalException(provider, error.Message, error); + } + else + { + var result = await read; + Assert.Equal(tableVersion, result.Version.Version); + this.fatalErrorHandler.DidNotReceiveWithAnyArgs().OnFatalException(default, default, default); + } + + Assert.Same(accepted, manager.MembershipTableSnapshot); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task DevelopmentMembershipReadChecksCommittedWriteVersion(bool readRow) + { + var cancellationToken = TestContext.Current.CancellationToken; + var backing = new InMemoryMembershipTable(new TableVersion(1, "1")); + var provider = CreateSystemTargetBasedMembershipTable(backing); + await provider.InitializeMembershipTableAsync(true, cancellationToken); + using var manager = CreateMembershipTableManager(provider); + await manager.UpdateStatus(SiloStatus.Joining, cancellationToken); + var accepted = manager.MembershipTableSnapshot; + Assert.Equal(new MembershipVersion(2), accepted.Version); + Assert.Equal(SiloStatus.Joining, manager.CurrentStatus); + backing.Version = new TableVersion(1, "reset"); + + var read = readRow + ? provider.ReadRowAsync(this.localSilo, cancellationToken) + : provider.ReadAllAsync(cancellationToken); + var error = await Assert.ThrowsAsync(() => read); + + Assert.Contains("version decreased from 2 to 1", error.Message, StringComparison.Ordinal); + this.fatalErrorHandler.Received(1).OnFatalException(provider, error.Message, error); + Assert.Same(accepted, manager.MembershipTableSnapshot); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task DevelopmentMembershipReadObservesCancellationBeforeReportingRollback(bool readRow) + { + var cancellationToken = TestContext.Current.CancellationToken; + var backing = Substitute.For(); + backing.ReadAllAsync(Arg.Any()) + .Returns(new MembershipTableData(new TableVersion(100, "100"))); + var provider = CreateSystemTargetBasedMembershipTable(backing); + await provider.InitializeMembershipTableAsync(true, cancellationToken); + using var manager = CreateMembershipTableManager(provider); + await manager.Refresh(cancellationToken: cancellationToken); + var accepted = manager.MembershipTableSnapshot; + var read = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + backing.ReadAllAsync(Arg.Any()).Returns(read.Task); + backing.ReadRowAsync(Arg.Any(), Arg.Any()).Returns(read.Task); + using var cancellation = new CancellationTokenSource(); + var pending = readRow + ? provider.ReadRowAsync(this.localSilo, cancellation.Token) + : provider.ReadAllAsync(cancellation.Token); + + cancellation.Cancel(); + read.SetResult(new MembershipTableData(new TableVersion(1, "1"))); + + var error = await Assert.ThrowsAnyAsync( + () => pending.WaitAsync(TimeSpan.FromSeconds(5), cancellationToken)); + Assert.Equal(cancellation.Token, error.CancellationToken); + Assert.Same(accepted, manager.MembershipTableSnapshot); + this.fatalErrorHandler.DidNotReceiveWithAnyArgs().OnFatalException(default, default, default); + } + + [Theory] + [InlineData(true)] + [InlineData(false)] + public async Task MembershipTableResetStopsOnlyTheVolatileProviderIncarnation(bool developmentProvider) + { + var cancellationToken = TestContext.Current.CancellationToken; + var original = await new InMemoryMembershipTable(new TableVersion(100, "old"), + Entry(this.localSilo, SiloStatus.Active, DateTimeOffset.UtcNow)).ReadAllAsync(cancellationToken); + var reset = await new InMemoryMembershipTable(new TableVersion(1, "new")).ReadAllAsync(cancellationToken); + var backing = new LegacyMembershipTable(Substitute.For()); + backing.ConfigureReadAll(original); + IMembershipTable provider = developmentProvider ? CreateSystemTargetBasedMembershipTable(backing) : backing; + await provider.InitializeMembershipTableAsync(true, cancellationToken); + using var manager = CreateMembershipTableManager(provider); + await manager.Refresh(cancellationToken: cancellationToken, requireFresh: true); + var accepted = manager.MembershipTableSnapshot; + Assert.Equal(new MembershipVersion(100), accepted.Version); + backing.ConfigureReadAll(reset); + + if (developmentProvider) + { + var error = await Assert.ThrowsAsync( + () => provider.ReadAllAsync(cancellationToken)); + Assert.Contains("version decreased from 100 to 1", error.Message, StringComparison.Ordinal); + await Assert.ThrowsAsync( + () => manager.Refresh(cancellationToken: cancellationToken, requireFresh: true)); + this.fatalErrorHandler.Received(1).OnFatalException(provider, error.Message, error); + } + else + { + await manager.Refresh(cancellationToken: cancellationToken, requireFresh: true); + this.fatalErrorHandler.DidNotReceiveWithAnyArgs().OnFatalException(default, default, default); + } + + Assert.Same(accepted, manager.MembershipTableSnapshot); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task OlderDevelopmentReadCompletingAfterNewerReadIsNotATableReset(bool readRow) + { + var cancellationToken = TestContext.Current.CancellationToken; + var original = await new InMemoryMembershipTable(new TableVersion(100, "old"), + Entry(this.localSilo, SiloStatus.Active, DateTimeOffset.UtcNow)).ReadAllAsync(cancellationToken); + var newer = new MembershipTableData(original.Members.ToList(), new TableVersion(101, "new")); + var backing = Substitute.For(); + backing.ReadAllAsync(Arg.Any()).Returns(original); + var provider = CreateSystemTargetBasedMembershipTable(backing); + await provider.InitializeMembershipTableAsync(true, cancellationToken); + using var manager = CreateMembershipTableManager(provider); + await manager.Refresh(cancellationToken: cancellationToken, requireFresh: true); + var read = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + backing.ReadAllAsync(Arg.Any()).Returns(read.Task); + backing.ReadRowAsync(Arg.Any(), Arg.Any()).Returns(read.Task); + Task pending = readRow + ? provider.ReadRowAsync(this.localSilo, cancellationToken) + : manager.Refresh(cancellationToken: cancellationToken, requireFresh: true); + try + { + Assert.False(pending.IsCompleted); + backing.ReadAllAsync(Arg.Any()).Returns(newer); + await manager.Refresh(cancellationToken: cancellationToken, requireFresh: true); + Assert.Equal(new MembershipVersion(101), manager.MembershipTableSnapshot.Version); + read.SetResult(original); + await pending.WaitAsync(TimeSpan.FromSeconds(5), cancellationToken); + + Assert.Equal(new MembershipVersion(101), manager.MembershipTableSnapshot.Version); + this.fatalErrorHandler.DidNotReceiveWithAnyArgs().OnFatalException(default, default, default); + } + finally + { + read.TrySetResult(original); + await pending.WaitAsync(TimeSpan.FromSeconds(5), cancellationToken); + } + } + + private sealed class MembershipContinuationContext : SynchronizationContext + { + private readonly Channel<(SendOrPostCallback Callback, object? State)> _continuations = + Channel.CreateUnbounded<(SendOrPostCallback, object?)>(); + + public override void Post(SendOrPostCallback callback, object? state) => + _continuations.Writer.TryWrite((callback, state)); + + public async Task<(SendOrPostCallback Callback, object? State)> TakeContinuation(CancellationToken cancellationToken) + { + using var deadline = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + deadline.CancelAfter(TimeSpan.FromSeconds(5)); + return await _continuations.Reader.ReadAsync(deadline.Token); + } } private sealed class BackoffTimeProvider : FakeTimeProvider @@ -1620,7 +1875,7 @@ private MembershipTableManager CreateMembershipTableManager( IAsyncTimerFactory? timerFactory = null, SiloLifecycleSubject? lifecycle = null) { - return new MembershipTableManager( + return this.membershipTableManager = new MembershipTableManager( localSiloDetails: this.localSiloDetails, clusterMembershipOptions: Options.Create(new ClusterMembershipOptions()), membershipTable: membershipTable, diff --git a/test/Orleans.Core.Tests/Membership/MembershipTableSnapshotTests.cs b/test/Orleans.Core.Tests/Membership/MembershipTableSnapshotTests.cs index 1aae925e5b4..f33ce06384e 100644 --- a/test/Orleans.Core.Tests/Membership/MembershipTableSnapshotTests.cs +++ b/test/Orleans.Core.Tests/Membership/MembershipTableSnapshotTests.cs @@ -14,6 +14,162 @@ namespace NonSilo.Tests.Membership [TestArea("Runtime")] public class MembershipTableSnapshotTests { + [Theory] + [InlineData("status")] + [InlineData("proxy-port")] + [InlineData("host")] + [InlineData("silo-name")] + [InlineData("role")] + [InlineData("update-zone")] + [InlineData("fault-zone")] + [InlineData("start-time")] + [InlineData("suspect-vote")] + public void SameVersionUpdatePreservesCanonicalFields(string field) + { + var local = Entry(Silo("127.0.0.1:100@1"), SiloStatus.Active, DateTimeOffset.UnixEpoch); + var dead = Entry(Silo("127.0.0.1:200@1"), SiloStatus.Dead, DateTimeOffset.UnixEpoch); + var previous = MembershipTableSnapshot.Create(Table(local, dead)); + var changed = local.WithIAmAliveTime(DateTime.UnixEpoch.AddMinutes(1)); + switch (field) + { + case "status": + changed.Status = SiloStatus.Dead; + break; + case "proxy-port": + changed.ProxyPort = 1234; + break; + case "host": + changed.HostName = "changed"; + break; + case "silo-name": + changed.SiloName = "changed"; + break; + case "role": + changed.RoleName = "changed"; + break; + case "update-zone": + changed.UpdateZone = 1; + break; + case "fault-zone": + changed.FaultZone = 1; + break; + case "start-time": + changed.StartTime = DateTime.UnixEpoch; + break; + case "suspect-vote": + changed.AddSuspector(dead.SiloAddress, DateTime.UnixEpoch); + break; + default: + throw new InvalidOperationException(field); + } + + var incoming = MembershipTableSnapshot.Create(Table(changed, dead)); + + Assert.All( + new[] { MembershipTableSnapshot.Update(previous, Table(changed, dead)), MembershipTableSnapshot.Update(previous, incoming) }, + updated => + { + Assert.Equal(2, updated.Entries.Count); + Assert.Contains(dead.SiloAddress, updated.Entries.Keys); + var retained = updated.Entries[local.SiloAddress]; + Assert.Equal(previous.Version, updated.Version); + Assert.True(updated.IsSuccessorTo(previous)); + Assert.Equal(local.SiloAddress, retained.SiloAddress); + Assert.Equal(local.Status, retained.Status); + Assert.Equal(local.ProxyPort, retained.ProxyPort); + Assert.Equal(local.HostName, retained.HostName); + Assert.Equal(local.SiloName, retained.SiloName); + Assert.Equal(local.RoleName, retained.RoleName); + Assert.Equal(local.UpdateZone, retained.UpdateZone); + Assert.Equal(local.FaultZone, retained.FaultZone); + Assert.Equal(local.StartTime, retained.StartTime); + Assert.Null(retained.SuspectTimes); + Assert.Equal(changed.IAmAliveTime, retained.IAmAliveTime); + }); + Assert.Equal(SiloStatus.Active, previous.Entries[local.SiloAddress].Status); + Assert.Equal(DateTime.UnixEpoch, previous.Entries[local.SiloAddress].IAmAliveTime); + Assert.Contains(dead.SiloAddress, previous.Entries.Keys); + } + + [Theory] + [InlineData(false, false)] + [InlineData(false, true)] + [InlineData(true, false)] + [InlineData(true, true)] + public void SameVersionUpdatePreservesCommittedFieldsAndMaximumHeartbeat(bool fromPeer, bool laterHeartbeat) + { + var local = Entry(Silo("127.0.0.1:100@1"), SiloStatus.Active, DateTimeOffset.UnixEpoch.AddMinutes(1)); + local.StartTime = DateTime.UnixEpoch.AddTicks(12345); + local.HostName = "committed-host"; + var previous = MembershipTableSnapshot.Create(Table(local)); + var persisted = local.WithIAmAliveTime(laterHeartbeat ? DateTime.UnixEpoch.AddMinutes(2) : DateTime.UnixEpoch); + persisted.StartTime = DateTime.UnixEpoch.AddMilliseconds(1); + persisted.HostName = "older-field"; + + var incoming = Table(persisted); + var updated = fromPeer + ? MembershipTableSnapshot.Update(previous, MembershipTableSnapshot.Create(incoming)) + : MembershipTableSnapshot.Update(previous, incoming); + + Assert.Equal(laterHeartbeat, updated.IsSuccessorTo(previous)); + var retained = Assert.Single(updated.Entries).Value; + Assert.Equal(local.HostName, retained.HostName); + Assert.Equal(local.StartTime, retained.StartTime); + Assert.Equal(laterHeartbeat ? persisted.IAmAliveTime : local.IAmAliveTime, retained.IAmAliveTime); + Assert.Equal("older-field", persisted.HostName); + } + + [Theory] + [InlineData(false, 1)] + [InlineData(false, 2)] + [InlineData(true, 1)] + [InlineData(true, 2)] + public void NewerVersionUpdateAdvancesCanonicalViewAndPreservesHeartbeat(bool fromPeer, int versionAdvance) + { + var local = Entry(Silo("127.0.0.1:100@1"), SiloStatus.Active, DateTimeOffset.UnixEpoch.AddMinutes(1)); + var previousTable = Table(local); + var previous = MembershipTableSnapshot.Create(previousTable); + var changed = local.WithStatus(SiloStatus.Dead).WithIAmAliveTime(DateTime.UnixEpoch); + var suspector = Silo("127.0.0.1:200@1"); + changed.AddSuspector(suspector, DateTime.UnixEpoch); + var incoming = new MembershipTableData( + new List> { Tuple.Create(changed, "updated") }, + new TableVersion(previousTable.Version.Version + versionAdvance, "updated")); + var updated = fromPeer + ? MembershipTableSnapshot.Update(previous, MembershipTableSnapshot.Create(incoming)) + : MembershipTableSnapshot.Update(previous, incoming); + + Assert.Equal(new MembershipVersion(previous.Version.Value + versionAdvance), updated.Version); + Assert.True(updated.IsSuccessorTo(previous)); + var retained = Assert.Single(updated.Entries).Value; + Assert.Equal(local.SiloAddress, retained.SiloAddress); + Assert.Equal(SiloStatus.Dead, retained.Status); + Assert.NotNull(retained.SuspectTimes); + var vote = Assert.Single(retained.SuspectTimes); + Assert.Equal(suspector, vote.Item1); + Assert.Equal(DateTime.UnixEpoch, vote.Item2); + Assert.Equal(local.IAmAliveTime, retained.IAmAliveTime); + Assert.Equal(SiloStatus.Active, previous.Entries[local.SiloAddress].Status); + Assert.Equal(DateTime.UnixEpoch, changed.IAmAliveTime); + } + + [Fact] + public void SameVersionHeartbeatAdvanceIsSuccessorBeforeStartTime() + { + var entry = Entry(Silo("127.0.0.1:100@1"), SiloStatus.Active, DateTimeOffset.UnixEpoch); + entry.StartTime = DateTime.UnixEpoch.AddMinutes(2); + var previous = MembershipTableSnapshot.Create(Table(entry)); + var incoming = MembershipTableSnapshot.Create(Table(entry.WithIAmAliveTime(DateTime.UnixEpoch.AddMinutes(1)))); + + var updated = MembershipTableSnapshot.Update(previous, incoming); + + Assert.True(updated.IsSuccessorTo(previous)); + var retained = Assert.Single(updated.Entries).Value; + Assert.Equal(entry.StartTime, retained.StartTime); + Assert.Equal(DateTime.UnixEpoch.AddMinutes(1), retained.IAmAliveTime); + Assert.Equal(DateTime.UnixEpoch, previous.Entries[entry.SiloAddress].IAmAliveTime); + } + [Fact] public void MembershipTableSnapshot_GetSiloStatus_JoiningSilo() { @@ -113,6 +269,132 @@ public void MembershipTableSnapshot_TryFormat_MatchesToString() AssertSpanFormattable(MembershipVersion.MinValue); } + [Theory] + [InlineData(SiloStatus.Created)] + [InlineData(SiloStatus.Joining)] + [InlineData(SiloStatus.Active)] + [InlineData(SiloStatus.ShuttingDown)] + [InlineData(SiloStatus.Stopping)] + public void SameVersionNonDeadRowRemovalIsNotASuccessor(SiloStatus removedStatus) + { + var active = Entry(Silo("127.0.0.1:100@1"), SiloStatus.Active, DateTimeOffset.UnixEpoch.AddMinutes(1)); + var removed = Entry(Silo("127.0.0.1:200@1"), removedStatus, DateTimeOffset.UnixEpoch); + var previous = MembershipTableSnapshot.Create(Table(active, removed)); + var incoming = MembershipTableSnapshot.Create(Table(active)); + var updated = MembershipTableSnapshot.Update(previous, incoming); + + Assert.Equal(previous.Version, updated.Version); + Assert.False(incoming.IsSuccessorTo(previous)); + Assert.False(updated.IsSuccessorTo(previous)); + Assert.False(previous.IsSuccessorTo(updated)); + Assert.Contains(removed.SiloAddress, previous.Entries.Keys); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public void SameVersionDeadRowPruningPreservesCanonicalFieldsAndHeartbeat(bool fromPeer) + { + var active = Entry(Silo("127.0.0.1:100@1"), SiloStatus.Active, DateTimeOffset.UnixEpoch.AddMinutes(2)); + active.HostName = "active-host"; + active.StartTime = DateTime.UnixEpoch; + var dead = Entry(Silo("127.0.0.1:200@1"), SiloStatus.Dead, DateTimeOffset.UnixEpoch); + var previous = MembershipTableSnapshot.Create(Table(active, dead)); + var retained = active.WithIAmAliveTime(DateTime.UnixEpoch.AddMinutes(1)); + var incoming = Table(retained); + + var updated = fromPeer + ? MembershipTableSnapshot.Update(previous, MembershipTableSnapshot.Create(incoming)) + : MembershipTableSnapshot.Update(previous, incoming); + + Assert.Equal(previous.Version, updated.Version); + Assert.True(updated.IsSuccessorTo(previous)); + Assert.False(previous.IsSuccessorTo(updated)); + Assert.False(updated.IsSuccessorTo(updated)); + var entry = Assert.Single(updated.Entries).Value; + Assert.Equal(active.SiloAddress, entry.SiloAddress); + Assert.Equal(SiloStatus.Active, entry.Status); + Assert.Equal(active.HostName, entry.HostName); + Assert.Equal(active.StartTime, entry.StartTime); + Assert.Equal(active.IAmAliveTime, entry.IAmAliveTime); + Assert.DoesNotContain(dead.SiloAddress, updated.Entries.Keys); + Assert.Contains(dead.SiloAddress, previous.Entries.Keys); + Assert.Equal(DateTime.UnixEpoch.AddMinutes(1), retained.IAmAliveTime); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public void VersionedDeadRowRemovalPreservesNewerLocalHeartbeat(bool fromPeer) + { + var address = Silo("127.0.0.1:100@1"); + var newer = DateTimeOffset.UnixEpoch.AddMinutes(2); + var older = DateTimeOffset.UnixEpoch.AddMinutes(1); + var previousTable = Table( + Entry(address, SiloStatus.Active, newer), + Entry(Silo("127.0.0.1:200@1"), SiloStatus.Dead)); + var previous = MembershipTableSnapshot.Create(previousTable); + var retained = Entry(address, SiloStatus.Active, older); + var incoming = new MembershipTableData(Tuple.Create(retained, "updated"), previousTable.Version.Next()); + var updated = fromPeer + ? MembershipTableSnapshot.Update(previous, MembershipTableSnapshot.Create(incoming)) + : MembershipTableSnapshot.Update(previous, incoming); + + Assert.Equal(new MembershipVersion(previous.Version.Value + 1), updated.Version); + Assert.True(updated.IsSuccessorTo(previous)); + Assert.Equal(newer.UtcDateTime, Assert.Single(updated.Entries).Value.IAmAliveTime); + Assert.Equal(older.UtcDateTime, retained.IAmAliveTime); + Assert.Equal(2, previous.Entries.Count); + } + + [Fact] + public void SameVersionRowReplacementIsNotASuccessor() + { + var active = Entry(Silo("127.0.0.1:100@1"), SiloStatus.Active); + var dead = Entry(Silo("127.0.0.1:200@1"), SiloStatus.Dead); + var previous = MembershipTableSnapshot.Create(Table(active, dead)); + var replacedEntry = MembershipTableSnapshot.Create(Table( + active.WithIAmAliveTime(DateTime.UnixEpoch.AddMinutes(1)), + Entry(Silo("127.0.0.1:300@1"), SiloStatus.Dead))); + + Assert.Equal(previous.Entries.Count, replacedEntry.Entries.Count); + Assert.False(replacedEntry.IsSuccessorTo(previous)); + Assert.Equal(SiloStatus.Active, previous.Entries[active.SiloAddress].Status); + } + + [Fact] + public void SameVersionDeadRowPruningCanRemoveAllRows() + { + var previous = MembershipTableSnapshot.Create(Table( + Entry(Silo("127.0.0.1:100@1"), SiloStatus.Dead), + Entry(Silo("127.0.0.1:200@1"), SiloStatus.Dead))); + Assert.All( + new[] { MembershipTableSnapshot.Update(previous, Table()), MembershipTableSnapshot.Update(previous, MembershipTableSnapshot.Create(Table())) }, + empty => + { + Assert.Equal(previous.Version, empty.Version); + Assert.True(empty.IsSuccessorTo(previous)); + Assert.False(empty.IsSuccessorTo(empty)); + Assert.Empty(empty.Entries); + }); + Assert.Equal(2, previous.Entries.Count); + } + + [Fact] + public void SameVersionHeartbeatAdvancePreservesNonDeadRows() + { + var keep = Entry(Silo("127.0.0.1:100@1"), SiloStatus.Active, DateTimeOffset.UnixEpoch); + var active = Entry(Silo("127.0.0.1:200@1"), SiloStatus.Active, DateTimeOffset.UnixEpoch); + var dead = Entry(Silo("127.0.0.1:300@1"), SiloStatus.Dead, DateTimeOffset.UnixEpoch); + var previous = MembershipTableSnapshot.Create(Table(keep, active, dead)); + var later = keep.WithIAmAliveTime(DateTime.UnixEpoch.AddMinutes(1)); + + Assert.False(MembershipTableSnapshot.Create(Table(later, dead)).IsSuccessorTo(previous)); + Assert.True(MembershipTableSnapshot.Create(Table(later, active)).IsSuccessorTo(previous)); + Assert.True(MembershipTableSnapshot.Create(Table(later, active, dead)).IsSuccessorTo(previous)); + Assert.Equal(DateTime.UnixEpoch, previous.Entries[keep.SiloAddress].IAmAliveTime); + } + private static SiloAddress Silo(string value) => SiloAddress.FromParsableString(value); private static MembershipEntry Entry(SiloAddress address, SiloStatus status, DateTimeOffset iAmAliveTime = default)