diff --git a/src/DotNext.Tests/Net/Cluster/Consensus/Raft/NetworkTransport/ConnectionOriented/Tcp/TcpTransportTests.cs b/src/DotNext.Tests/Net/Cluster/Consensus/Raft/NetworkTransport/ConnectionOriented/Tcp/TcpTransportTests.cs index 39f9f8bd4..f125e62f2 100644 --- a/src/DotNext.Tests/Net/Cluster/Consensus/Raft/NetworkTransport/ConnectionOriented/Tcp/TcpTransportTests.cs +++ b/src/DotNext.Tests/Net/Cluster/Consensus/Raft/NetworkTransport/ConnectionOriented/Tcp/TcpTransportTests.cs @@ -96,6 +96,28 @@ public Task StressTest() return StressTestCore(CreateServer, CreateClient); } + [Fact] + public Task RequestTimeout() + { + static TcpServer CreateServer(ILocalMember member, EndPoint address, TimeSpan timeout) => new(address, 2, member, NullLoggerFactory.Instance) + { + MemoryAllocator = MemoryAllocator.Default, + ReceiveTimeout = timeout, + TransmissionBlockSize = 65535, + GracefulShutdownTimeout = 2000 + }; + + static TcpClient CreateClient(EndPoint address, ILocalMember member, TimeSpan timeout) => new(member, address) + { + MemoryAllocator = MemoryAllocator.Default, + RequestTimeout = timeout, + ConnectTimeout = timeout, + TransmissionBlockSize = 65535, + }; + + return RequestTimeoutTest(CreateServer, CreateClient); + } + [Theory] [InlineData(true)] [InlineData(false)] @@ -393,4 +415,4 @@ private static RaftCluster.TcpConfiguration CreateConfiguration(int port, bool c return result; } -} \ No newline at end of file +} diff --git a/src/DotNext.Tests/Net/Cluster/Consensus/Raft/NetworkTransport/TransportTestSuite.cs b/src/DotNext.Tests/Net/Cluster/Consensus/Raft/NetworkTransport/TransportTestSuite.cs index a8cdb0034..7c30cdb54 100644 --- a/src/DotNext.Tests/Net/Cluster/Consensus/Raft/NetworkTransport/TransportTestSuite.cs +++ b/src/DotNext.Tests/Net/Cluster/Consensus/Raft/NetworkTransport/TransportTestSuite.cs @@ -38,6 +38,7 @@ private sealed class LocalMember : Assert, ILocalMember internal ReceiveEntriesBehavior Behavior; internal byte[] ReceivedConfiguration = []; internal long ReceivedConfigurationVersion = -1L; + internal TimeSpan VoteDelay; private readonly ClusterMemberId localId = Random.Shared.Next(); internal LocalMember(bool smallAmountOfMetadata = false) @@ -125,13 +126,17 @@ async ValueTask ILocalMember.InstallConfigurationAsync(lon return true; } - ValueTask> ILocalMember.VoteAsync(ClusterMemberId sender, long term, long lastLogIndex, long lastLogTerm, CancellationToken token) + async ValueTask> ILocalMember.VoteAsync(ClusterMemberId sender, long term, long lastLogIndex, long lastLogTerm, CancellationToken token) { True(token.CanBeCanceled); Equal(42L, term); Equal(1L, lastLogIndex); Equal(56L, lastLogTerm); - return ValueTask.FromResult>(new() { Term = 43L, Value = true }); + + if (VoteDelay > TimeSpan.Zero) + await Task.Delay(VoteDelay, token); + + return new() { Term = 43L, Value = true }; } ValueTask> ILocalMember.PreVoteAsync(ClusterMemberId sender, long term, long lastLogIndex, long lastLogTerm, CancellationToken token) @@ -212,6 +217,27 @@ private protected async Task StressTestCore(ServerFactory serverFactory, ClientF }); } + private protected async Task RequestTimeoutTest(ServerFactory serverFactory, ClientFactory clientFactory) + { + var serverAddr = new IPEndPoint(IPAddress.Loopback, 3789); + var serverTimeout = TimeSpan.FromMilliseconds(50D); + var member = new LocalMember { VoteDelay = DefaultTimeout }; + await using var server = serverFactory(member, serverAddr, serverTimeout); + await server.StartAsync(TestToken); + + using (var client = clientFactory(serverAddr, member, DefaultTimeout)) + { + await ThrowsAsync( + () => client.As().VoteAsync(42L, 1L, 56L, TestToken)); + } + + member.VoteDelay = TimeSpan.Zero; + using var nextClient = clientFactory(serverAddr, member, DefaultTimeout); + var result = await nextClient.As().VoteAsync(42L, 1L, 56L, TestToken); + True(result.Value); + Equal(43L, result.Term); + } + private protected async Task MetadataRequestResponseTest(ServerFactory serverFactory, ClientFactory clientFactory, bool smallAmountOfMetadata) { var timeout = DefaultTimeout; @@ -494,4 +520,4 @@ private protected async Task LeadershipCore(Func new(new() { Location = GetTempPath() }, IStateMachine.CreateNoOp()); -} \ No newline at end of file +} diff --git a/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/NetworkTransport/ConnectionOriented/Tcp/TcpServer.cs b/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/NetworkTransport/ConnectionOriented/Tcp/TcpServer.cs index 104230fb9..c048ed7c4 100644 --- a/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/NetworkTransport/ConnectionOriented/Tcp/TcpServer.cs +++ b/src/cluster/DotNext.Net.Cluster/Net/Cluster/Consensus/Raft/NetworkTransport/ConnectionOriented/Tcp/TcpServer.cs @@ -131,6 +131,11 @@ private async void HandleConnection(Socket remoteClient) { await ProcessRequestAsync(messageType, protocol, timeoutSource.Token).ConfigureAwait(false); } + catch (OperationCanceledException e) when (e.CausedByTimeout(timeoutSource)) + { + logger.RequestTimedOut(clientAddress, e); + break; + } finally { // reset cancellation token @@ -142,11 +147,9 @@ private async void HandleConnection(Socket remoteClient) { logger.ConnectionWasResetByClient(clientAddress); } - catch (OperationCanceledException e) + catch (OperationCanceledException) { - // if lifecycleToken is canceled then shutdown socket gracefully without logging - if (e.CausedByTimeout(timeoutSource)) - logger.RequestTimedOut(clientAddress, e); + // shutdown socket gracefully without logging } catch (Exception e) { @@ -267,4 +270,4 @@ protected override async ValueTask DisposeAsyncCore() logger.TcpGracefulShutdownFailed(GracefulShutdownTimeout); } } -} \ No newline at end of file +}