Hacktoberfest 2026:メンテナが10月に向けて印を付けた、オープンで初心者向けの issue。 Hacktoberfest の issue を見る

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

オープン 初心者向け
#301 コメント 0 件 リアクション 0 件 担当者 0 名 GitHub で見る

まだ誰も着手していません。

評価

難易度
2/5
見積もり時間
1〜3時間
初心者へのやさしさ
78/100
issue の種類
バグ
明瞭さ
明確に書かれている
活発さ
活発
技術スタック
csharp

調査の方向性

RaftCluster.Configuration.cs の CustomTransportConfiguration.CreateClient から始め、GenericClient の初期化を BuiltInTransportConfiguration.CreateClient と比較します。Repro.csproj と Program.cs を使用して 30 秒のタイムアウトを再現し、その後、設定した 500 ms の RequestTimeout がカスタムトランスポートのクライアントに適用され、サイレントメンバーへのリクエストがその時間に近い時間で失敗することを確認します。

索引モデルが issue の本文から書いたものです。

説明

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")].

主要言語
C#
スター
2k
フォーク
159
平均マージ
1日 12時間
マージ済み PR(30日)
1

環境構築

Codespaces で開く

このプロジェクトの開発コンテナを、あなたの GitHub アカウントでブラウザ上に起動します。

はじめの一歩

  1. issue を最後まで読み、次にプロジェクトのコントリビューションガイドを読みます。
  2. 着手することを issue にコメントします — 二人が同じ作業をするのを防げます。
  3. リポジトリをフォークし、ブランチを切って変更します。
  4. issue 番号を参照したプルリクエストを送ります。

dotnet/dotNext のほかの issue

dotnet/dotNext の issue をすべて見る

似ている issue

C# の issue をもっと見る

新しい issue をメールで受け取る

初心者向けの GitHub issue を短くまとめたダイジェスト。