From fab605be2a52cf8a1a238a0a9f4d6b01ccecae64 Mon Sep 17 00:00:00 2001 From: Guillaume Chervet Date: Tue, 15 Sep 2026 17:39:57 +0200 Subject: [PATCH 1/3] Fix Raft member rejoining and preserve existing WAL metadata pages Signed-off-by: Guillaume Chervet --- CHANGELOG.md | 6 ++ .../Raft/Http/RaftHttpClusterTests.cs | 50 ++++++++++++++++ .../Raft/StateMachine/WriteAheadLogTests.cs | 57 +++++++++++++++++++ .../DotNext.Net.Cluster/ExceptionMessages.cs | 4 +- .../ExceptionMessages.restext | 3 +- .../Net/Cluster/Consensus/Raft/RaftCluster.cs | 12 +++- .../WriteAheadLog.MetadataManagement.cs | 23 +++++++- .../WriteAheadLog.PageManagement.cs | 6 +- .../Raft/StateMachine/WriteAheadLog.cs | 21 +++---- 9 files changed, 164 insertions(+), 18 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 8cec404ec..2e9b67beb 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,12 @@ 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. + # 09-11-2026 DotNext 6.7.2 * Minor performance improvements of static extension methods declared in `AdvancedHelpers` class 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/StateMachine/WriteAheadLogTests.cs b/src/DotNext.Tests/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLogTests.cs index 4a1aef3ec..f7bcb04d0 100644 --- a/src/DotNext.Tests/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLogTests.cs +++ b/src/DotNext.Tests/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLogTests.cs @@ -17,6 +17,63 @@ namespace DotNext.Net.Cluster.Consensus.Raft.StateMachine; [Collection(TestCollections.WriteAheadLog)] public sealed class WriteAheadLogTests : Test { + [Theory] + [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 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], TestToken); + 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 }, TestToken)); + await wal.CommitAsync(count, TestToken); + await wal.WaitForApplyAsync(count, TestToken); + await wal.FlushAsync(TestToken); + } + + All(metadata.EnumerateFiles(), file => Equal(pageSize, file.Length)); + await using var reopened = new WriteAheadLog(options, new ContextAwareStateMachine()); + await reopened.InitializeAsync(TestToken); + await reopened.ReadAsync(new LogEntryConsumer(async (entries, _, token) => + { + 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: token)); + } + return Missing.Value; + }), 1L, count, TestToken); + } + + [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() { 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 +} From 750c57fdae57d6d195df5e92dc97c526c08cbfff Mon Sep 17 00:00:00 2001 From: Guillaume Chervet Date: Tue, 15 Sep 2026 17:51:18 +0200 Subject: [PATCH 2/3] Preserve Raft HTTP compatibility with pre-versioned peers Signed-off-by: Guillaume Chervet --- CHANGELOG.md | 5 +- src/DotNext.Tests/DotNext.Tests.csproj | 3 +- .../Raft/Http/RaftHttpMessageTests.cs | 76 +++++++++++++++++++ .../Raft/StateMachine/WriteAheadLogTests.cs | 4 +- .../Cluster/Messaging/MessageHandlerTests.cs | 10 ++- .../Cluster/Messaging/TestMessageHandler.cs | 4 +- .../DotNext.AspNetCore.Cluster/Assembly.cs | 7 +- .../Raft/Http/AppendEntriesMessage.cs | 9 ++- .../Consensus/Raft/Http/RaftHttpMessage.cs | 8 +- 9 files changed, 113 insertions(+), 13 deletions(-) create mode 100644 src/DotNext.Tests/Net/Cluster/Consensus/Raft/Http/RaftHttpMessageTests.cs diff --git a/CHANGELOG.md b/CHANGELOG.md index 2e9b67beb..b22a64c5d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,9 @@ 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 @@ -3453,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/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/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 f7bcb04d0..6392c41dd 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; @@ -503,7 +505,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 +} From 91572acc15f44a1791a203072fd9f47e74ae7113 Mon Sep 17 00:00:00 2001 From: Guillaume Chervet Date: Tue, 15 Sep 2026 18:09:24 +0200 Subject: [PATCH 3/3] Bound WAL regression tests and report macOS test progress Signed-off-by: Guillaume Chervet --- azure-pipelines.yml | 4 ++-- .../Raft/StateMachine/WriteAheadLogTests.cs | 21 ++++++++++--------- 2 files changed, 13 insertions(+), 12 deletions(-) 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/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLogTests.cs b/src/DotNext.Tests/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLogTests.cs index 6392c41dd..8cb631a68 100644 --- a/src/DotNext.Tests/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLogTests.cs +++ b/src/DotNext.Tests/Net/Cluster/Consensus/Raft/StateMachine/WriteAheadLogTests.cs @@ -19,42 +19,43 @@ namespace DotNext.Net.Cluster.Consensus.Raft.StateMachine; [Collection(TestCollections.WriteAheadLog)] public sealed class WriteAheadLogTests : Test { - [Theory] + [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], TestToken); + 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 }, TestToken)); - await wal.CommitAsync(count, TestToken); - await wal.WaitForApplyAsync(count, TestToken); - await wal.FlushAsync(TestToken); + 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(TestToken); - await reopened.ReadAsync(new LogEntryConsumer(async (entries, _, token) => + 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: token)); + Equal($"entry-{i + 1}", await entries[i].ToStringAsync(Encoding.UTF8, token: readToken)); } return Missing.Value; - }), 1L, count, TestToken); + }), 1L, count, token); } [Theory]