Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,8 @@ The client and silo must use matching gateway ports, <xref:Orleans.Configuration
- Configure a shared development clustering primary and assign each silo unique silo and gateway ports.
- Run the same production clustering provider against a local container or emulator.

Development clustering keeps its membership table in the primary silo's memory, with membership versions increasing throughout the cluster's lifetime. Losing that table invalidates the cluster; restart all of its silos to establish a new cluster. A surviving silo which reads a table version below the membership version it already knows invokes <xref:Orleans.Runtime.IFatalErrorHandler.OnFatalException*>. 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]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 <xref:Orleans.Runtime.SiloStatus> 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 <xref:Orleans.Runtime.SiloStatus> 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 <xref:Orleans.IMembershipTable.UpdateIAmAlive*> 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 <xref:Orleans.IMembershipTable.UpdateIAmAlive*> 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
Expand Down
73 changes: 45 additions & 28 deletions src/Orleans.Core/Runtime/MembershipTableSnapshot.cs
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,12 @@
namespace Orleans.Runtime
{
/// <summary>
/// Represents an immutable snapshot of cluster membership state.
/// Represents an immutable snapshot of a canonical membership view and per-silo liveness timestamps.
/// </summary>
/// <remarks>
/// 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.
/// </remarks>
[GenerateSerializer, Immutable]
internal sealed class MembershipTableSnapshot : ISpanFormattable
{
Expand Down Expand Up @@ -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<MembershipEntry> updatedEntries)
private static MembershipTableSnapshot Update(
MembershipTableSnapshot previousSnapshot,
MembershipVersion version,
IEnumerable<MembershipEntry> updatedEntries)
{
ArgumentNullException.ThrowIfNull(previousSnapshot);
ArgumentNullException.ThrowIfNull(updatedEntries);

var entries = ImmutableDictionary.CreateBuilder<SiloAddress, MembershipEntry>();
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());
}

/// <summary>
Expand Down Expand Up @@ -150,6 +156,10 @@ public SiloStatus GetSiloStatus(SiloAddress silo)
/// <summary>
/// Determines whether this snapshot is a successor to another snapshot.
/// </summary>
/// <remarks>
/// At the same canonical membership version, progress consists of newer liveness timestamps
/// or pruning previously Dead rows. Non-Dead rows are retained.
/// </remarks>
/// <param name="other">The snapshot to compare against.</param>
/// <returns><see langword="true"/> if this snapshot is a successor to <paramref name="other"/>; otherwise, <see langword="false"/>.</returns>
public bool IsSuccessorTo(MembershipTableSnapshot other)
Expand All @@ -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()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -186,6 +186,7 @@ public async Task RefreshFromSnapshot(MembershipTableSnapshot snapshot, Cancella

LogInformationReceivedClusterMembershipSnapshot(this.log, snapshot);

cancellationToken.ThrowIfCancellationRequested();
this.TryProcessMembershipUpdate(MembershipTableSnapshot.Update, snapshot, nameof(RefreshFromSnapshot));
}

Expand Down Expand Up @@ -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)
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<SystemTargetBasedMembershipTable> logger)
public SystemTargetBasedMembershipTable(
IServiceProvider serviceProvider,
ILogger<SystemTargetBasedMembershipTable> 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);

Expand Down Expand Up @@ -103,12 +111,52 @@ private async Task WaitForTableGrainToInit(IMembershipTableSystemTarget membersh
[Obsolete("Use ReadRowAsync instead.")]
public Task<MembershipTableData> ReadRow(SiloAddress key) => ReadRowAsync(key, CancellationToken.None);

public Task<MembershipTableData> ReadRowAsync(SiloAddress key, CancellationToken cancellationToken = default) => this.grain.ReadRowAsync(key, cancellationToken);
public async Task<MembershipTableData> 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<MembershipTableData> ReadAll() => ReadAllAsync(CancellationToken.None);

public Task<MembershipTableData> ReadAllAsync(CancellationToken cancellationToken = default) => this.grain.ReadAllAsync(cancellationToken);
public async Task<MembershipTableData> 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<IMembershipManager>();
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<bool> InsertRow(MembershipEntry entry, TableVersion tableVersion) => InsertRowAsync(entry, tableVersion, CancellationToken.None);
Expand Down
Loading
Loading