Hacktoberfest 2026: le issue che i maintainer hanno segnato per ottobre, aperte e adatte ai principianti. Sfoglia le issue Hacktoberfest

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

Aperta Adatta ai principianti
#301 0 commenti 0 reazioni 0 assegnatari Vedi su GitHub

Nessuno ha ancora preso questa issue.

Valutazione

Difficoltà
2/5
Tempo stimato
1-3 ore
Idoneità per principianti
78/100
Tipo di issue
Bug
Chiarezza
Specificata chiaramente
Stato di attività
Attiva
Stack tecnologico
csharp

Direzione di ricerca

Inizia in RaftCluster.Configuration.cs, in CustomTransportConfiguration.CreateClient, e confronta la sua inizializzazione di GenericClient con BuiltInTransportConfiguration.CreateClient. Usa Repro.csproj e Program.cs per riprodurre il timeout di 30 secondi, quindi verifica che un RequestTimeout configurato a 500 ms venga applicato ai client con trasporto personalizzato e che la richiesta al membro silenzioso fallisca dopo un intervallo prossimo a quella durata.

Scritto dal modello di indicizzazione a partire dal testo della issue.

Descrizione

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

Lingua principale
C#
Stelle
2k
Fork
160
Merge medio
1g 12h
PR unite (30g)
1

Preparare l'ambiente

Come iniziare

  1. Leggi tutta la issue e poi la guida ai contributi del progetto.
  2. Commenta sulla issue per dire che te ne occupi tu — evita che due persone facciano lo stesso lavoro.
  3. Fai un fork del repository e lavora su un branch.
  4. Apri una pull request che faccia riferimento al numero della issue.

Altre issue di dotnet/dotNext

Tutte le issue di dotnet/dotNext

Issue simili

Altre issue su C#

Ricevi le nuove issue nella tua casella

Un breve riepilogo di issue GitHub adatte ai principianti.