From 20c9bd00d2115ac3feadab410dc1b623b0780f31 Mon Sep 17 00:00:00 2001 From: wlodzimierzyk Date: Fri, 14 Aug 2026 17:40:06 +0200 Subject: [PATCH 1/2] fix tcp request timeout handling --- .../Tcp/TcpTransportTests.cs | 24 +++++++++++++- .../NetworkTransport/TransportTestSuite.cs | 32 +++++++++++++++++-- .../ConnectionOriented/Tcp/TcpServer.cs | 13 +++++--- 3 files changed, 60 insertions(+), 9 deletions(-) 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 +} From cbe0c9805ef242337a8ebf1b8d9c204bbd0058e4 Mon Sep 17 00:00:00 2001 From: wlodzimierzyk Date: Mon, 17 Aug 2026 08:28:37 +0200 Subject: [PATCH 2/2] address tcp timeout review feedback --- .../ConnectionOriented/Tcp/TcpServer.cs | 26 +++++++------------ 1 file changed, 10 insertions(+), 16 deletions(-) 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 c048ed7c4..d34efc253 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 @@ -127,29 +127,22 @@ private async void HandleConnection(Socket remoteClient) break; timeoutSource = multiplexer.Combine(receiveTimeout, lifecycleToken); - try - { - await ProcessRequestAsync(messageType, protocol, timeoutSource.Token).ConfigureAwait(false); - } - catch (OperationCanceledException e) when (e.CausedByTimeout(timeoutSource)) - { - logger.RequestTimedOut(clientAddress, e); - break; - } - finally - { - // reset cancellation token - await timeoutSource.DisposeAsync().ConfigureAwait(false); - } + await ProcessRequestAsync(messageType, protocol, timeoutSource.Token).ConfigureAwait(false); + + // reset cancellation token + await timeoutSource.DisposeAsync().ConfigureAwait(false); + timeoutSource = default; } } catch (Exception e) when (e is SocketException { SocketErrorCode: SocketError.ConnectionReset } or { InnerException: SocketException { SocketErrorCode: SocketError.ConnectionReset } }) { logger.ConnectionWasResetByClient(clientAddress); } - catch (OperationCanceledException) + catch (OperationCanceledException e) { - // shutdown socket gracefully without logging + // if lifecycleToken is canceled then shutdown socket gracefully without logging + if (e.CausedByTimeout(timeoutSource)) + logger.RequestTimedOut(clientAddress, e); } catch (Exception e) { @@ -157,6 +150,7 @@ private async void HandleConnection(Socket remoteClient) } finally { + await timeoutSource.DisposeAsync().ConfigureAwait(false); await protocol.DisposeAsync().ConfigureAwait(false); if (protocol.BaseStream is SslStream ssl) await ssl.DisposeAsync().ConfigureAwait(false);