diff --git a/CHANGELOG.md b/CHANGELOG.md
index 8cec404ec..b22a64c5d 100644
--- a/CHANGELOG.md
+++ b/CHANGELOG.md
@@ -1,6 +1,15 @@
Release Notes
====
+# Unreleased
+
+DotNext.Net.Cluster:
+* Preserve the size of existing WAL metadata pages when reopening logs from older releases or hosts with a different system page size. Reject inconsistent page sizes before opening WAL files.
+* Allow a removed live member to receive replication and rejoin after its election waiters have faulted. Restore normal follower behavior and fresh leadership waiters on re-addition.
+
+DotNext.AspNetCore.Cluster:
+* Accept legacy Raft HTTP requests without a state machine version as version zero and responses without the last-index backtracking hint. Malformed explicit headers remain rejected, allowing rolling upgrades from 6.6.0 without relaxing version checks.
+
# 09-11-2026
DotNext 6.7.2
* Minor performance improvements of static extension methods declared in `AdvancedHelpers` class
@@ -3447,4 +3456,4 @@ This release introduces a new feature called Value Delegates which are allocatio
DotNext.Net.Cluster 0.2.0
DotNext.AspNetCore.Cluster 0.2.0
-* Raft client is now capable to ensure that changes are committed by leader node using [WriteConcern](https://dotnet.github.io/dotNext/versions/1.x/api/DotNext.Net.Cluster.Replication.WriteConcern.html)
\ No newline at end of file
+* Raft client is now capable to ensure that changes are committed by leader node using [WriteConcern](https://dotnet.github.io/dotNext/versions/1.x/api/DotNext.Net.Cluster.Replication.WriteConcern.html)
diff --git a/azure-pipelines.yml b/azure-pipelines.yml
index c7a514f95..4810a9b1a 100644
--- a/azure-pipelines.yml
+++ b/azure-pipelines.yml
@@ -132,7 +132,7 @@ stages:
inputs:
command: test
projects: $(TestProject)
- arguments: --configuration Debug --coverage --coverage-output-format cobertura --coverage-output $(Agent.TempDirectory)/CoverageResults/coverage.cobertura.xml
+ arguments: --configuration Debug --output Detailed --coverage --coverage-output-format cobertura --coverage-output $(Agent.TempDirectory)/CoverageResults/coverage.cobertura.xml
nobuild: false
testRunTitle: 'Debug on MacOS'
publishTestResults: true
@@ -146,7 +146,7 @@ stages:
inputs:
command: test
projects: $(AotTestProject)
- arguments: --configuration Debug --coverage --coverage-output-format cobertura --coverage-output $(Agent.TempDirectory)/CoverageResults/coverage.cobertura.xml
+ arguments: --configuration Debug --output Detailed --coverage --coverage-output-format cobertura --coverage-output $(Agent.TempDirectory)/CoverageResults/coverage.cobertura.xml
nobuild: false
testRunTitle: 'Debug on MacOS (NoDynamicCode)'
publishTestResults: true
diff --git a/src/DotNext.Tests/DotNext.Tests.csproj b/src/DotNext.Tests/DotNext.Tests.csproj
index ffb68afaf..2d96b0fe5 100644
--- a/src/DotNext.Tests/DotNext.Tests.csproj
+++ b/src/DotNext.Tests/DotNext.Tests.csproj
@@ -30,7 +30,8 @@
-
+
+
diff --git a/src/DotNext.Tests/Net/Cluster/Consensus/Raft/Http/RaftHttpClusterTests.cs b/src/DotNext.Tests/Net/Cluster/Consensus/Raft/Http/RaftHttpClusterTests.cs
index 8f9c73d10..5e1c062fa 100644
--- a/src/DotNext.Tests/Net/Cluster/Consensus/Raft/Http/RaftHttpClusterTests.cs
+++ b/src/DotNext.Tests/Net/Cluster/Consensus/Raft/Http/RaftHttpClusterTests.cs
@@ -51,6 +51,56 @@ private static IHost CreateHost(int port, IDictionary
.Build();
}
+ [Fact(Timeout = 60000)]
+ public static async Task RemovedLiveMemberCanRejoin()
+ {
+ var token = TestContext.Current.CancellationToken;
+ static Dictionary Configuration(int port, bool coldStart) => new()
+ {
+ ["partitioning"] = "false",
+ ["publicEndPoint"] = $"http://localhost:{port}",
+ ["coldStart"] = coldStart.ToString(),
+ ["requestTimeout"] = "00:00:05",
+ };
+
+ using var host1 = CreateHost(3262, Configuration(3262, true));
+ using var host2 = CreateHost(3263, Configuration(3263, false));
+ using var host3 = CreateHost(3264, Configuration(3264, false));
+ await host1.StartAsync(token);
+ await host2.StartAsync(token);
+ await host3.StartAsync(token);
+ var leader = GetLocalClusterView(host1);
+ var second = GetLocalClusterView(host2);
+ var removed = GetLocalClusterView(host3);
+ await leader.WaitForLeaderAsync(DefaultTimeout, token);
+ True(await leader.AddMemberAsync(second.LocalMemberAddress, token));
+ True(await leader.AddMemberAsync(removed.LocalMemberAddress, token));
+ await removed.Readiness.WaitAsync(token);
+
+ True(await leader.RemoveMemberAsync(removed.LocalMemberAddress, token));
+ await leader.ReplicateAsync(new EmptyLogEntry { Term = leader.Term }, token);
+ True(await leader.AddMemberAsync(removed.LocalMemberAddress, token));
+ await leader.ForceReplicationAsync(token);
+ var index = leader.AuditTrail.LastCommittedEntryIndex;
+ // Catch-up applies the removal before receiving the re-addition. The old
+ // election task has faulted, but subsequent AppendEntries must still work.
+ await removed.AuditTrail.WaitForApplyAsync(index, token).AsTask().WaitAsync(DefaultTimeout, token);
+ Equal(leader.LocalMemberAddress, ((UriEndPoint)removed.Leader.EndPoint).Uri);
+ await leader.ReplicateAsync(new EmptyLogEntry { Term = leader.Term }, token);
+ await removed.AuditTrail.WaitForApplyAsync(leader.AuditTrail.LastCommittedEntryIndex, token)
+ .AsTask().WaitAsync(DefaultTimeout, token);
+
+ using var leadershipWaitCancellation = CancellationTokenSource.CreateLinkedTokenSource(token);
+ var leadershipWait = removed.WaitForLeadershipAsync(leadershipWaitCancellation.Token);
+ False(leadershipWait.IsCompleted);
+ await leadershipWaitCancellation.CancelAsync();
+ await ThrowsAnyAsync(leadershipWait);
+
+ await host3.StopAsync(token);
+ await host2.StopAsync(token);
+ await host1.StopAsync(token);
+ }
+
[Fact]
public static async Task CommunicationWithLeader()
{
diff --git a/src/DotNext.Tests/Net/Cluster/Consensus/Raft/Http/RaftHttpMessageTests.cs b/src/DotNext.Tests/Net/Cluster/Consensus/Raft/Http/RaftHttpMessageTests.cs
new file mode 100644
index 000000000..48f2eeb4e
--- /dev/null
+++ b/src/DotNext.Tests/Net/Cluster/Consensus/Raft/Http/RaftHttpMessageTests.cs
@@ -0,0 +1,76 @@
+using System.Net.Http;
+using Microsoft.AspNetCore.Http;
+
+namespace DotNext.Net.Cluster.Consensus.Raft.Http;
+
+public sealed class RaftHttpMessageTests : Test
+{
+ private static HttpRequest CreateVoteRequest(string stateVersion)
+ {
+ var message = new RequestVoteMessage(default, term: 2L, lastLogIndex: 5L, lastLogTerm: 1L, stateVersion: 0);
+ using var outgoing = new HttpRequestMessage();
+ message.PrepareRequest(outgoing);
+ var request = new DefaultHttpContext().Request;
+ foreach (var header in outgoing.Headers)
+ request.Headers[header.Key] = header.Value.ToArray();
+
+ if (stateVersion is null)
+ request.Headers.Remove("X-Raft-State-Version");
+ else
+ request.Headers["X-Raft-State-Version"] = stateVersion;
+
+ return request;
+ }
+
+ [Theory]
+ [InlineData(null, 0)]
+ [InlineData("0", 0)]
+ [InlineData("42", 42)]
+ public static void LegacyStateVersionDefaultsToZero(string header, int expected)
+ {
+ var message = new RequestVoteMessage(CreateVoteRequest(header));
+ Equal(expected, message.StateVersion);
+ Equal(2L, message.ConsensusTerm);
+ Equal(5L, message.LastLogIndex);
+ }
+
+ [Fact]
+ public static void MalformedStateVersionIsRejected()
+ => Throws(() => new RequestVoteMessage(CreateVoteRequest("invalid")));
+
+ private static IHttpMessage> CreateAppendRequest()
+ => new AppendEntriesMessage(default, term: 2L,
+ prevLogIndex: 5L, prevLogTerm: 1L, commitIndex: 4L, entries: [], stateVersion: 0);
+
+ private static HttpResponseMessage CreateAppendResponse(HeartbeatResult result, string lastIndex)
+ {
+ var response = new HttpResponseMessage { Content = new StringContent(result.ToString()) };
+ response.Headers.Add("X-Raft-Term", "2");
+ if (lastIndex is not null)
+ response.Headers.Add("X-Raft-Last-Index", lastIndex);
+
+ return response;
+ }
+
+ [Theory]
+ [InlineData(HeartbeatResult.Rejected, null, 5L)]
+ [InlineData(HeartbeatResult.ReplicatedWithLeaderTerm, null, 5L)]
+ [InlineData(HeartbeatResult.Rejected, "3", 3L)]
+ [InlineData(HeartbeatResult.ReplicatedWithLeaderTerm, "7", 7L)]
+ public static async Task LegacyAppendResponseFallsBackToPreviousIndex(HeartbeatResult result, string header, long expected)
+ {
+ using var response = CreateAppendResponse(result, header);
+ var parsed = await CreateAppendRequest().ParseResponseAsync(response, TestContext.Current.CancellationToken);
+ Equal(2L, parsed.Term);
+ Equal(result, parsed.Value.Result);
+ Equal(expected, parsed.Value.LastIndex);
+ }
+
+ [Fact]
+ public static async Task MalformedLastIndexIsRejected()
+ {
+ using var response = CreateAppendResponse(HeartbeatResult.Rejected, "invalid");
+ await ThrowsAsync(() => CreateAppendRequest()
+ .ParseResponseAsync(response, TestContext.Current.CancellationToken));
+ }
+}
diff --git a/src/DotNext.Tests/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLogTests.cs b/src/DotNext.Tests/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLogTests.cs
index 4a1aef3ec..8cb631a68 100644
--- a/src/DotNext.Tests/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLogTests.cs
+++ b/src/DotNext.Tests/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLogTests.cs
@@ -1,3 +1,5 @@
+extern alias RaftCore;
+using CoreConfigurationStorage = RaftCore::DotNext.Net.Cluster.Consensus.Raft.Membership.InMemoryClusterConfigurationStorage;
using System.Buffers.Binary;
using System.Net;
using System.Reflection;
@@ -17,6 +19,64 @@ namespace DotNext.Net.Cluster.Consensus.Raft.StateMachine;
[Collection(TestCollections.WriteAheadLog)]
public sealed class WriteAheadLogTests : Test
{
+ [Theory(Timeout = 60000)]
+ [InlineData(4096, WriteAheadLog.MemoryManagementStrategy.SharedMemory)]
+ [InlineData(16384, WriteAheadLog.MemoryManagementStrategy.SharedMemory)]
+ [InlineData(4096, WriteAheadLog.MemoryManagementStrategy.PrivateMemory)]
+ [InlineData(16384, WriteAheadLog.MemoryManagementStrategy.PrivateMemory)]
+ public static async Task ExistingMetadataPageSizeIsPreserved(int pageSize, WriteAheadLog.MemoryManagementStrategy strategy)
+ {
+ var token = TestContext.Current.CancellationToken;
+ var directory = GetTempPath();
+ var metadata = Directory.CreateDirectory(Path.Combine(directory, "metadata"));
+ // A valid empty page, laid out by either the 4 KiB legacy format or a
+ // 16 KiB-page host. This also exercises cross-host reopening on 4 KiB CI.
+ await File.WriteAllBytesAsync(Path.Combine(metadata.FullName, "0"), new byte[pageSize], token);
+ var options = new WriteAheadLog.Options { Location = directory, MemoryManagement = strategy };
+ const int count = 1025;
+ await using (var wal = new WriteAheadLog(options, new ContextAwareStateMachine()))
+ {
+ for (var i = 1; i <= count; i++)
+ Equal(i, await wal.AppendAsync(new TestLogEntry($"entry-{i}") { Term = i }, token));
+ await wal.CommitAsync(count, token);
+ await wal.WaitForApplyAsync(count, token);
+ await wal.FlushAsync(token);
+ }
+
+ All(metadata.EnumerateFiles(), file => Equal(pageSize, file.Length));
+ await using var reopened = new WriteAheadLog(options, new ContextAwareStateMachine());
+ await reopened.InitializeAsync(token);
+ await reopened.ReadAsync(new LogEntryConsumer(async (entries, _, readToken) =>
+ {
+ Equal(count, entries.Count);
+ for (var i = 0; i < entries.Count; i++)
+ {
+ Equal(i + 1L, entries[i].Term);
+ Equal($"entry-{i + 1}", await entries[i].ToStringAsync(Encoding.UTF8, token: readToken));
+ }
+ return Missing.Value;
+ }), 1L, count, token);
+ }
+
+ [Theory]
+ [InlineData(0)]
+ [InlineData(4097)]
+ [InlineData(8192)] // Individually valid, but inconsistent with the other page.
+ public static void InvalidMetadataPageSizeIsRejectedBeforeOpeningWal(int secondPageSize)
+ {
+ var directory = GetTempPath();
+ var metadata = Directory.CreateDirectory(Path.Combine(directory, "metadata"));
+ var first = Path.Combine(metadata.FullName, "0");
+ var second = Path.Combine(metadata.FullName, "1");
+ File.WriteAllBytes(first, new byte[4096]);
+ File.WriteAllBytes(second, new byte[secondPageSize]);
+ Throws(() => new WriteAheadLog(new() { Location = directory }, new ContextAwareStateMachine()));
+ Equal(4096L, new FileInfo(first).Length);
+ Equal(secondPageSize, new FileInfo(second).Length);
+ False(File.Exists(Path.Combine(directory, "checkpoint")));
+ False(File.Exists(Path.Combine(directory, "state")));
+ }
+
[Fact]
public static async Task LockManager()
{
@@ -446,7 +506,7 @@ public static async Task CaptureConfiguration()
{
var dir = GetTempPath();
await using var wal = new WriteAheadLog(new() { Location = dir }, IStateMachine.CreateNoOp(2));
- IClusterConfigurationStorage storage = new InMemoryClusterConfigurationStorage(EqualityComparer.Default);
+ IClusterConfigurationStorage storage = new CoreConfigurationStorage(EqualityComparer.Default);
wal.ConfigurationStorage = storage;
var config = await storage.LoadConfigurationAsync(TestToken);
diff --git a/src/DotNext.Tests/Net/Cluster/Messaging/MessageHandlerTests.cs b/src/DotNext.Tests/Net/Cluster/Messaging/MessageHandlerTests.cs
index 030c5e834..37c6848cd 100644
--- a/src/DotNext.Tests/Net/Cluster/Messaging/MessageHandlerTests.cs
+++ b/src/DotNext.Tests/Net/Cluster/Messaging/MessageHandlerTests.cs
@@ -1,3 +1,5 @@
+extern alias RaftCore;
+using CoreMessageHandler = RaftCore::DotNext.Net.Cluster.Messaging.MessageHandler;
namespace DotNext.Net.Cluster.Messaging;
public sealed class MessageHandlerTests : Test
@@ -5,7 +7,7 @@ public sealed class MessageHandlerTests : Test
[Fact]
public static void MessageHandlerBuilder1()
{
- var handler = new MessageHandler.Builder()
+ var handler = new CoreMessageHandler.Builder()
.Add(AddMessage.Name, static (sender, input, context, token) => Task.FromResult(input.Execute()), ResultMessage.Name)
.Add(ResultMessage.Name, static (sender, input, context, token) => Task.CompletedTask)
.Build();
@@ -21,7 +23,7 @@ public static void MessageHandlerBuilder1()
[Fact]
public static void MessageHandlerBuilder2()
{
- var handler = new MessageHandler.Builder()
+ var handler = new CoreMessageHandler.Builder()
.Add(AddMessage.Name, static (input, context, token) => Task.FromResult(input.Execute()), ResultMessage.Name)
.Add(ResultMessage.Name, static (ResultMessage input, object context, CancellationToken token) => Task.CompletedTask)
.Build();
@@ -37,7 +39,7 @@ public static void MessageHandlerBuilder2()
[Fact]
public static void MessageHandlerBuilder3()
{
- var handler = new MessageHandler.Builder()
+ var handler = new CoreMessageHandler.Builder()
.Add(AddMessage.Name, static (sender, input, token) => Task.FromResult(input.Execute()), ResultMessage.Name)
.Add(ResultMessage.Name, static (ISubscriber sender, ResultMessage input, CancellationToken token) => Task.CompletedTask)
.Build();
@@ -53,7 +55,7 @@ public static void MessageHandlerBuilder3()
[Fact]
public static void MessageHandlerBuilder4()
{
- var handler = new MessageHandler.Builder()
+ var handler = new CoreMessageHandler.Builder()
.Add(AddMessage.Name, static (input, token) => Task.FromResult(input.Execute()), ResultMessage.Name)
.Add(ResultMessage.Name, static (input, token) => Task.CompletedTask)
.Build();
diff --git a/src/DotNext.Tests/Net/Cluster/Messaging/TestMessageHandler.cs b/src/DotNext.Tests/Net/Cluster/Messaging/TestMessageHandler.cs
index 8bc35b32c..cf3afa519 100644
--- a/src/DotNext.Tests/Net/Cluster/Messaging/TestMessageHandler.cs
+++ b/src/DotNext.Tests/Net/Cluster/Messaging/TestMessageHandler.cs
@@ -1,3 +1,5 @@
+extern alias RaftCore;
+using CoreMessageHandler = RaftCore::DotNext.Net.Cluster.Messaging.MessageHandler;
using System.Diagnostics.CodeAnalysis;
@@ -7,7 +9,7 @@ namespace DotNext.Net.Cluster.Messaging;
[Message(AddMessage.Name)]
[Message(SubtractMessage.Name)]
[Message(ResultMessage.Name)]
-public class TestMessageHandler : MessageHandler
+public class TestMessageHandler : CoreMessageHandler
{
internal int Result;
diff --git a/src/cluster/DotNext.AspNetCore.Cluster/Assembly.cs b/src/cluster/DotNext.AspNetCore.Cluster/Assembly.cs
index df5e65922..9c5ac9f68 100644
--- a/src/cluster/DotNext.AspNetCore.Cluster/Assembly.cs
+++ b/src/cluster/DotNext.AspNetCore.Cluster/Assembly.cs
@@ -1,4 +1,9 @@
using System.Runtime.InteropServices;
+#if DEBUG
+using System.Runtime.CompilerServices;
+
+[assembly: InternalsVisibleTo("DotNext.Tests")]
+#endif
[assembly: CLSCompliant(true)]
-[assembly: ComVisible(false)]
\ No newline at end of file
+[assembly: ComVisible(false)]
diff --git a/src/cluster/DotNext.AspNetCore.Cluster/Net/Cluster/Consensus/Raft/Http/AppendEntriesMessage.cs b/src/cluster/DotNext.AspNetCore.Cluster/Net/Cluster/Consensus/Raft/Http/AppendEntriesMessage.cs
index d17e582ef..65783be09 100644
--- a/src/cluster/DotNext.AspNetCore.Cluster/Net/Cluster/Consensus/Raft/Http/AppendEntriesMessage.cs
+++ b/src/cluster/DotNext.AspNetCore.Cluster/Net/Cluster/Consensus/Raft/Http/AppendEntriesMessage.cs
@@ -489,9 +489,14 @@ async Task> IHttpMessage>.Pa
Term = ParseTerm(response),
Value = new()
{
- LastIndex = ParseHeader(response.Headers, LastIndexHeader, Int64Parser),
+ // Older peers do not provide this backtracking hint. Falling back
+ // to PrevLogIndex preserves the legacy one-entry decrement on rejection;
+ // successful replication uses the last index sent by the leader.
+ LastIndex = response.Headers.Contains(LastIndexHeader)
+ ? ParseHeader(response.Headers, LastIndexHeader, Int64Parser)
+ : PrevLogIndex,
Result = await HttpMessage.ParseEnumResponseAsync(response, token).ConfigureAwait(false),
}
};
}
-}
\ No newline at end of file
+}
diff --git a/src/cluster/DotNext.AspNetCore.Cluster/Net/Cluster/Consensus/Raft/Http/RaftHttpMessage.cs b/src/cluster/DotNext.AspNetCore.Cluster/Net/Cluster/Consensus/Raft/Http/RaftHttpMessage.cs
index e902786fb..5bbe9b0cc 100644
--- a/src/cluster/DotNext.AspNetCore.Cluster/Net/Cluster/Consensus/Raft/Http/RaftHttpMessage.cs
+++ b/src/cluster/DotNext.AspNetCore.Cluster/Net/Cluster/Consensus/Raft/Http/RaftHttpMessage.cs
@@ -28,7 +28,11 @@ private protected RaftHttpMessage(IDictionary headers)
: base(headers)
{
ConsensusTerm = ParseHeader(headers, TermHeader, Int64Parser);
- StateVersion = ParseHeader(headers, StateVersionHeader, Int32Parser);
+ // Peers predating state machine versioning implicitly use version zero.
+ // A present but malformed version must still fail protocol validation.
+ StateVersion = headers.ContainsKey(StateVersionHeader)
+ ? ParseHeader(headers, StateVersionHeader, Int32Parser)
+ : 0;
}
protected new void PrepareRequest(HttpRequestMessage request)
@@ -72,4 +76,4 @@ private protected static Task SaveResponseAsync(HttpResponse response, in Res
WriteTerm(response, result.Term);
return SaveResponseAsync(response, result.Value, token);
}
-}
\ No newline at end of file
+}
diff --git a/src/cluster/DotNext.Net.Cluster/ExceptionMessages.cs b/src/cluster/DotNext.Net.Cluster/ExceptionMessages.cs
index 872e83f16..2fbd1b603 100644
--- a/src/cluster/DotNext.Net.Cluster/ExceptionMessages.cs
+++ b/src/cluster/DotNext.Net.Cluster/ExceptionMessages.cs
@@ -49,9 +49,11 @@ internal static string UnknownRaftMessageType(T messageType)
internal static string BadCheckpointVersion(uint version) => Resources.Get().Format(version);
+ internal static string InvalidWalMetadataPageSize => (string)Resources.Get();
+
internal static string LogEntryHashMismatch => (string)Resources.Get();
internal static string MissingWalPage(uint pageIndex) => Resources.Get().Format(pageIndex);
internal static string StateMachineIsNotRestored => (string)Resources.Get();
-}
\ No newline at end of file
+}
diff --git a/src/cluster/DotNext.Net.Cluster/ExceptionMessages.restext b/src/cluster/DotNext.Net.Cluster/ExceptionMessages.restext
index e0e3cb29b..f12ad369c 100644
--- a/src/cluster/DotNext.Net.Cluster/ExceptionMessages.restext
+++ b/src/cluster/DotNext.Net.Cluster/ExceptionMessages.restext
@@ -18,4 +18,5 @@ BadProtocolVersion=Multiplexing protocol version is not supported: {0}
BadCheckpointVersion=Checkpoint file version is not supported: {0}
LogEntryHashMismatch=Log entry hash doesn't match
MissingWalPage=WAL page {0} doesn't exist on the disk
-StateMachineIsNotRestored=State machine is not restored. Call RestoreAsync first.
\ No newline at end of file
+StateMachineIsNotRestored=State machine is not restored. Call RestoreAsync first.
+InvalidWalMetadataPageSize=WAL metadata pages must have the same power-of-two size of at least 4096 bytes
diff --git a/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/RaftCluster.cs b/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/RaftCluster.cs
index 782b129ef..9df943e7d 100644
--- a/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/RaftCluster.cs
+++ b/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/RaftCluster.cs
@@ -252,7 +252,8 @@ private set
Debug.Assert(value is not null);
raiseEventHandlers = electionEventCopy.TrySetResult(value);
break;
- case (true, false) when !ReferenceEquals(electionEventCopy.Task.Result, value):
+ case (true, false) when !electionEventCopy.Task.IsCompletedSuccessfully
+ || !ReferenceEquals(electionEventCopy.Task.Result, value):
Debug.Assert(value is not null);
newEvent = new();
newEvent.SetResult(value);
@@ -298,7 +299,7 @@ private ValueTask UnfreezeAsync()
ValueTask result;
// ensure that local member has been received
- if (readinessProbe.Task.IsCompleted)
+ if (readinessProbe.Task.IsCompleted && state is not StandbyState { Resumable: false })
{
result = ValueTask.CompletedTask;
}
@@ -316,6 +317,11 @@ private ValueTask UnfreezeAsync()
async ValueTask UnfreezeCoreAsync()
{
+ // A removed member can be added again after its original readiness
+ // probe completed and its leadership waiters were faulted.
+ if (leadershipEvent.Task.IsFaulted)
+ Interlocked.Exchange(ref leadershipEvent, new(TaskCreationOptions.RunContinuationsAsynchronously));
+
var newState = new FollowerState(this) { ConsensusReached = true };
await UpdateStateAsync(newState).ConfigureAwait(false);
newState.StartServing(ElectionTimeout);
@@ -1471,4 +1477,4 @@ file static class LockTypes
public sealed class TransitionLock : AsyncExclusiveLock;
public sealed class MembershipLock : AsyncExclusiveLock;
-}
\ No newline at end of file
+}
diff --git a/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLog.MetadataManagement.cs b/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLog.MetadataManagement.cs
index d35d376e2..190416cd6 100644
--- a/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLog.MetadataManagement.cs
+++ b/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLog.MetadataManagement.cs
@@ -17,6 +17,27 @@ private readonly struct MetadataPageManager(PageManager manager, int hashSizeInB
public const string LocationPrefix = "metadata";
private readonly int MetadataEntryAlignedSize = GetAlignedSize(LogEntryMetadata.Size + hashSizeInBytes, manager.PageSize);
+
+ internal static int GetPageSize(DirectoryInfo location)
+ {
+ int? storedSize = null;
+ foreach (var file in location.EnumerateFiles())
+ {
+ if (!uint.TryParse(file.Name, provider: null, out _))
+ continue;
+
+ var length = file.Length;
+ if (length is < Page.MinSize or > int.MaxValue || !long.IsPow2(length)
+ || (storedSize.HasValue && storedSize.GetValueOrDefault() != length))
+ throw new InvalidDataException(ExceptionMessages.InvalidWalMetadataPageSize);
+
+ storedSize = (int)length;
+ }
+
+ // Older WALs use 4 KiB metadata pages even on hosts with larger OS
+ // pages. Their page numbering must survive both upgrades and moves.
+ return storedSize ?? int.Max(Page.MinSize, Environment.SystemPageSize);
+ }
private static int GetAlignedSize(int headerSize, int containerSize)
{
@@ -147,4 +168,4 @@ static MemoryManager IMetadataView.GetPage(PageManager man
static MetadataWriter IMetadataView.Create(Span buffer)
=> new(buffer);
}
-}
\ No newline at end of file
+}
diff --git a/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLog.PageManagement.cs b/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLog.PageManagement.cs
index de16551eb..0a061aee8 100644
--- a/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLog.PageManagement.cs
+++ b/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLog.PageManagement.cs
@@ -390,7 +390,9 @@ static nuint GetAlignment(int pageSize, out nint madvise)
// fallback - no THP/LP support
madvise = 0;
- alignment = (uint)Environment.SystemPageSize;
+ // Legacy metadata pages may be smaller than an OS page. They
+ // cannot use discard/huge-page hints, but remain valid buffers.
+ alignment = (uint)int.Min(pageSize, Environment.SystemPageSize);
exit:
return alignment;
@@ -467,4 +469,4 @@ protected override MemoryMappedPage CreatePage(uint pageIndex)
protected override void ReleasePage(MemoryMappedPage page)
=> page.As().Dispose();
}
-}
\ No newline at end of file
+}
diff --git a/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLog.cs b/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLog.cs
index 0706a425f..37a8561e1 100644
--- a/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLog.cs
+++ b/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLog.cs
@@ -50,11 +50,15 @@ public WriteAheadLog(Options configuration, IStateMachine stateMachine)
// Snapshot getter may throw if the state machine is not restored or initialized
var snapshotIndex = stateMachine.Snapshot?.Index ?? 0L;
+ var rootPath = new DirectoryInfo(configuration.Location);
+ rootPath.CreateIfNeeded();
+ var metadataLocation = rootPath.GetSubdirectory(MetadataPageManager.LocationPrefix);
+ metadataLocation.CreateIfNeeded();
+ // Validate existing pages before opening or resizing any WAL files.
+ var metadataPageSize = MetadataPageManager.GetPageSize(metadataLocation);
hash = configuration.CreateHashAlgorithm();
lifetimeToken = (lifetimeTokenSource = new()).Token;
cancellationTokens = new();
- var rootPath = new DirectoryInfo(configuration.Location);
- rootPath.CreateIfNeeded();
context = new(DictionaryConcurrencyLevel, configuration.ConcurrencyLevel);
lockManager = new()
@@ -89,9 +93,6 @@ public WriteAheadLog(Options configuration, IStateMachine stateMachine)
// page management
{
- var metadataLocation = rootPath.GetSubdirectory(MetadataPageManager.LocationPrefix);
- metadataLocation.CreateIfNeeded();
-
var dataLocation = rootPath.GetSubdirectory(PagedBufferWriter.LocationPrefix);
dataLocation.CreateIfNeeded();
@@ -99,20 +100,20 @@ public WriteAheadLog(Options configuration, IStateMachine stateMachine)
switch (configuration.MemoryManagement)
{
case MemoryManagementStrategy.PrivateMemory when OperatingSystem.IsWindows() && configuration.NoBuffering:
- m = new WindowsDirectPageManager(metadataLocation, int.Max(Page.MinSize, Environment.SystemPageSize));
+ m = new WindowsDirectPageManager(metadataLocation, metadataPageSize);
d = new WindowsDirectPageManager(dataLocation, configuration.ChunkSize);
break;
case MemoryManagementStrategy.PrivateMemory when OperatingSystem.IsLinux() && configuration.NoBuffering:
- m = new LinuxDirectPageManager(metadataLocation, int.Max(Page.MinSize, Environment.SystemPageSize));
+ m = new LinuxDirectPageManager(metadataLocation, metadataPageSize);
d = new LinuxDirectPageManager(dataLocation, configuration.ChunkSize);
break;
case MemoryManagementStrategy.PrivateMemory:
- m = new AnonymousPageManager(metadataLocation, int.Max(Page.MinSize, Environment.SystemPageSize));
+ m = new AnonymousPageManager(metadataLocation, metadataPageSize);
d = new AnonymousPageManager(dataLocation, configuration.ChunkSize);
break;
case MemoryManagementStrategy.SharedMemory:
default:
- m = new MemoryMappedPageManager(metadataLocation, int.Max(Page.MinSize, Environment.SystemPageSize));
+ m = new MemoryMappedPageManager(metadataLocation, metadataPageSize);
d = new MemoryMappedPageManager(dataLocation, configuration.ChunkSize);
break;
}
@@ -695,4 +696,4 @@ public static void CreateIfNeeded(this DirectoryInfo directory)
public static DirectoryInfo GetSubdirectory(this DirectoryInfo root, string prefix)
=> new(Path.Combine(root.FullName, prefix));
-}
\ No newline at end of file
+}