Hacktoberfest 2026: những issue maintainer đã đánh dấu cho tháng Mười, đang mở và phù hợp người mới. Xem issue Hacktoberfest

CustomTransportConfiguration ignores RequestTimeout for its clients (members keep 30 s)

Đang mở Phù hợp với người mới
#301 0 bình luận 0 reaction 0 người được giao Xem trên GitHub

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ả

ai_assisted Lib:Cluster

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

Mở trong Codespaces

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.

Bắt đầu từ đâu

  1. Đọc hết issue, rồi đọc hướng dẫn đóng góp của dự án.
  2. 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.
  3. Fork repository và làm thay đổi trên một nhánh.
  4. Mở pull request có tham chiếu số hiệu của issue.

Issue khác của dotnet/dotNext

Tất cả issue của dotnet/dotNext

Issue tương tự

Thêm issue về C#

Nhận issue mới trong hộp thư của bạn

Bản tóm tắt ngắn những issue GitHub phù hợp với người mới.