CustomTransportConfiguration ignores RequestTimeout for its clients (members keep 30 s)
Chưa có ai nhận issue này.
Đánh giá
- Độ khó
- 2/5
- Thời gian dự kiến
- 1-3 giờ
- Mức phù hợp với người mới
- 78/100
- Loại issue
- Lỗi
- Độ rõ ràng
- Đặc tả rõ ràng
- Mức độ hoạt động
- Sôi nổi
- Công nghệ
- csharp
- Lĩnh vực
- distributed-systems
Hướng nghiên cứu
Bắt đầu trong RaftCluster.Configuration.cs tại CustomTransportConfiguration.CreateClient và so sánh việc khởi tạo GenericClient của nó với BuiltInTransportConfiguration.CreateClient. Sử dụng Repro.csproj và Program.cs để tái hiện timeout 30 giây, sau đó xác minh rằng RequestTimeout được cấu hình ở mức 500 ms được áp dụng cho các client sử dụng custom transport và yêu cầu tới thành viên im lặng thất bại sau khoảng thời gian gần bằng thời lượng đó.
Do mô hình lập chỉ mục viết ra từ nội dung của issue.
Mô tả
Package: DotNext.Net.Cluster 6.8.1 (.NET 10)
Summary
RaftCluster.CustomTransportConfiguration.CreateClient creates GenericClient without setting
RequestTimeout, so every member on a custom transport keeps RaftClusterMember's default of 30 s,
whatever RequestTimeout is configured. BuiltInTransportConfiguration.CreateClient (TCP) sets it:
// RaftCluster.Configuration.cs, CustomTransportConfiguration
internal override GenericClient CreateClient(ILocalMember localMember, EndPoint endPoint) => new(localMember, endPoint)
{
DefaultAllocator = MemoryAllocator,
ConnectTimeout = ConnectTimeout,
ConnectionFactory = clientFactory,
// RequestTimeout = RequestTimeout is missing (the TCP transport sets it)
};
Why it matters
A request to a member that stops answering (a partitioned or black-holed peer, which doesn't refuse
the connection) then waits 30 s. In a 3-member cluster with one member gone silent, we intermittently
saw the leader's heartbeat rounds stop reaching a majority: the healthy follower heard nothing, started
elections it couldn't win, and the leader's lease lapsed. A dump of a stuck leader showed DoHeartbeats
waiting on its replication barrier and a member's Client.RequestAsync awaiting an AppendEntries
response with a request timeout of 30 s (not the configured 500 ms).
Our reading of why the healthy follower stalls too (from the source, not proven by the dump): the
AppendEntries exchange holds the write-ahead log's read lock while the entries are sent, so the next
append waits for the write lock behind it, and the other members' replication reads queue behind that
append. Either way, with RequestTimeout applied (we now set it on every member ourselves, see the
workaround below), the stall no longer reproduces in our tests.
Reproduction
Two nodes on the custom transport; node A's connections to B can be silenced (bytes held both ways).
B joins, then goes silent, and A asks it for its metadata:
Configured RequestTimeout: 500 ms
B answering: request took 7 ms
B silent: request failed with MemberUnavailableException after 30005 ms (expected about 500 ms)
Repro.csproj:
<Project Sdk="Microsoft.NET.Sdk.Web">
<PropertyGroup>
<OutputType>Exe</OutputType>
<TargetFramework>net10.0</TargetFramework>
<Nullable>enable</Nullable>
<ImplicitUsings>enable</ImplicitUsings>
<ManagePackageVersionsCentrally>false</ManagePackageVersionsCentrally>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="DotNext.Net.Cluster" Version="6.8.1" />
</ItemGroup>
</Project>
Program.cs:
// Repro: RaftCluster.CustomTransportConfiguration ignores RequestTimeout for its clients
// (DotNext.Net.Cluster 6.8.1). Requests to a peer that stops answering wait 30 s (RaftClusterMember's
// default request timeout), not the configured RequestTimeout.
//
// Two nodes on the custom transport. A's connections to B can be silenced (bytes held both ways, as a
// network partition does). B joins normally; then it goes silent and A asks it for its metadata.
using System.Diagnostics;
using System.IO.Pipelines;
using System.Net;
using System.Net.Sockets;
using DotNext.Net.Cluster.Consensus.Raft;
using DotNext.Net.Cluster.Consensus.Raft.Membership;
using DotNext.Net.Cluster.Consensus.Raft.StateMachine;
using Microsoft.AspNetCore.Connections;
using Microsoft.AspNetCore.Http.Features;
using Microsoft.AspNetCore.Server.Kestrel.Transport.Sockets;
using Microsoft.Extensions.Logging.Abstractions;
using Microsoft.Extensions.Options;
var requestTimeout = TimeSpan.FromMilliseconds(500);
var a = await Node.Start(coldStart: true, requestTimeout);
var b = await Node.Start(coldStart: false, requestTimeout);
await a.Cluster.WaitForLeaderAsync(TimeSpan.FromSeconds(10));
if (!await a.Cluster.AddMemberAsync(b.EndPoint, CancellationToken.None))
{
throw new InvalidOperationException("B was not added.");
}
var member = a.Cluster.Members.First(m => m.EndPoint.Equals(b.EndPoint));
Console.WriteLine($"Configured RequestTimeout: {requestTimeout.TotalMilliseconds} ms");
var watch = Stopwatch.StartNew();
await member.GetMetadataAsync(refresh: true, CancellationToken.None);
Console.WriteLine($"B answering: request took {watch.Elapsed.TotalMilliseconds:F0} ms");
Gate.Silenced = true;
watch.Restart();
try
{
await member.GetMetadataAsync(refresh: true, CancellationToken.None);
Console.WriteLine("B silent: request succeeded?");
}
catch (Exception e)
{
Console.WriteLine($"B silent: request failed with {e.GetType().Name} after {watch.Elapsed.TotalMilliseconds:F0} ms (expected about {requestTimeout.TotalMilliseconds} ms)");
}
Environment.Exit(0);
sealed class Node
{
public required RaftCluster Cluster { get; init; }
public required IPEndPoint EndPoint { get; init; }
public static async Task<Node> Start(bool coldStart, TimeSpan requestTimeout)
{
var endPoint = new IPEndPoint(IPAddress.Loopback, FreePort());
var directory = Directory.CreateTempSubdirectory("dotnext-repro-").FullName;
var configuration = new RaftCluster.CustomTransportConfiguration(
endPoint,
new SocketTransportFactory(Options.Create(new SocketTransportOptions()), NullLoggerFactory.Instance),
new GatedConnections())
{
LowerElectionTimeout = 1000,
UpperElectionTimeout = 2000,
RequestTimeout = requestTimeout,
ColdStart = coldStart,
ConfigurationStorage = new EndPointStorage(Path.Combine(directory, "members")),
};
var wal = new WriteAheadLog(new() { Location = Path.Combine(directory, "wal") }, IStateMachine.CreateNoOp());
var cluster = new RaftCluster(configuration) { AuditTrail = wal };
await cluster.StartAsync(CancellationToken.None);
return new Node { Cluster = cluster, EndPoint = endPoint };
}
private static int FreePort()
{
var listener = new TcpListener(IPAddress.Loopback, 0);
listener.Start();
var port = ((IPEndPoint)listener.LocalEndpoint).Port;
listener.Stop();
return port;
}
}
static class Gate
{
public static volatile bool Silenced;
}
/// <summary>Outbound TCP connections whose bytes are held, both ways, while <see cref="Gate.Silenced"/>.</summary>
sealed class GatedConnections : IConnectionFactory
{
private readonly SocketConnectionContextFactory _sockets = new(new SocketConnectionFactoryOptions(), NullLogger.Instance);
public async ValueTask<ConnectionContext> ConnectAsync(EndPoint endpoint, CancellationToken cancellationToken = default)
{
var socket = new Socket(endpoint.AddressFamily, SocketType.Stream, ProtocolType.Tcp);
await socket.ConnectAsync(endpoint, cancellationToken);
return new GatedConnection(_sockets.Create(socket));
}
}
sealed class GatedConnection : ConnectionContext
{
private readonly ConnectionContext _inner;
private readonly Pipe _input = new();
private readonly Pipe _output = new();
public GatedConnection(ConnectionContext inner)
{
_inner = inner;
Transport = new DuplexPipe(_input.Reader, _output.Writer);
_ = Pump(_inner.Transport.Input, _input.Writer);
_ = Pump(_output.Reader, _inner.Transport.Output);
}
public override IDuplexPipe Transport { get; set; }
public override string ConnectionId { get => _inner.ConnectionId; set => _inner.ConnectionId = value; }
public override IFeatureCollection Features => _inner.Features;
public override IDictionary<object, object?> Items { get => _inner.Items; set => _inner.Items = value; }
public override EndPoint? RemoteEndPoint { get => _inner.RemoteEndPoint; set => _inner.RemoteEndPoint = value; }
public override ValueTask DisposeAsync() => _inner.DisposeAsync();
private static async Task Pump(PipeReader from, PipeWriter to)
{
try
{
while (true)
{
var read = await from.ReadAsync();
while (Gate.Silenced)
{
await Task.Delay(50);
}
foreach (var segment in read.Buffer)
{
await to.WriteAsync(segment);
}
from.AdvanceTo(read.Buffer.End);
if (read.IsCompleted)
{
break;
}
}
}
catch
{
}
}
private sealed class DuplexPipe(PipeReader input, PipeWriter output) : IDuplexPipe
{
public PipeReader Input => input;
public PipeWriter Output => output;
}
}
/// <summary>Keeps the Raft configuration: endpoints as text.</summary>
sealed class EndPointStorage(string path) : PersistentClusterConfigurationStorage<EndPoint>(path)
{
protected override void Encode(EndPoint address, ref DotNext.Buffers.BufferWriterSlim<byte> writer)
{
var bytes = System.Text.Encoding.UTF8.GetBytes(address.ToString()!);
Span<byte> length = stackalloc byte[4];
System.Buffers.Binary.BinaryPrimitives.WriteInt32LittleEndian(length, bytes.Length);
writer.Write(length);
writer.Write(bytes);
}
protected override EndPoint Decode(ref DotNext.Buffers.SequenceReader reader)
{
var bytes = new byte[reader.ReadLittleEndian<int>()];
reader.Read(bytes);
return IPEndPoint.Parse(System.Text.Encoding.UTF8.GetString(bytes));
}
}
Expected
CustomTransportConfiguration.CreateClient sets RequestTimeout = RequestTimeout, as the TCP
transport does.
Workaround
Setting RaftClusterMember's private requestTimeout field on every member (after start, and whenever
a peer is discovered), e.g. with [UnsafeAccessor(UnsafeAccessorKind.Field, Name = "requestTimeout")].
- Ngôn ngữ chính
- C#
- Star
- 2k
- Fork
- 159
- Merge trung bình
- 1 ngày 12 giờ
- Pull request đã merge (30 ngày)
- 1
Chuẩn bị môi trường
Khởi chạy dev container của dự án ngay trên trình duyệt, bằng tài khoản GitHub của bạn.
- Không có Dockerfile hay tệp Docker Compose
- Không có mẫu pull request
- Đọc hướng dẫn đóng góp
Bắt đầu từ đâu
- Đọc hết issue, rồi đọc hướng dẫn đóng góp của dự án.
- Bình luận trên issue rằng bạn sẽ nhận — tránh hai người làm cùng một việc.
- Fork repository và làm thay đổi trên một nhánh.
- Mở pull request có tham chiếu số hiệu của issue.
Issue khác của dotnet/dotNext
-
Lib:Threading question wontfix
Độ khó 4/5 3-5 ngày Mức phù hợp với người mới 45/100
Tất cả issue của dotnet/dotNext
Issue tương tự
-
Money ExploitsĐang mởS: Untriaged
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 62/100
project-wayfarer/wayfarer-14#1628 ·
Maintainer thường phản hồi trong vòng 3 ngày
-
:watch: Not Triaged dotnet-target-version
Độ khó 1/5 Dưới một giờ Mức phù hợp với người mới 85/100
Maintainer thường phản hồi trong vòng 1 ngày
-
copilot documentation
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 88/100
Maintainer thường phản hồi trong vòng 2 ngày
-
untriaged
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 86/100
dotnet/dotnet-api-docs#13124 ·
Maintainer thường phản hồi trong vòng 1 ngày
-
agentic-workflows
Độ khó 2/5 1-3 giờ Mức phù hợp với người mới 74/100
Maintainer thường phản hồi trong vòng 1 ngày