diff --git a/docker-compose.yml b/docker-compose.yml index 0896f11148..5f2865ed01 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -16,7 +16,7 @@ services: env_file: - shared.env environment: - - EVENTSTORE_GOSSIP_SEED=172.30.240.12:2113,172.30.240.13:2113 + - EVENTSTORE_GOSSIP_SEED=172.30.240.12:1112,172.30.240.13:1112 - EVENTSTORE_REPLICATION_IP=172.30.240.11 - EVENTSTORE_CERTIFICATE_FILE=/etc/eventstore/certs/node/node.crt - EVENTSTORE_CERTIFICATE_PRIVATE_KEY_FILE=/etc/eventstore/certs/node/node.key @@ -43,7 +43,7 @@ services: env_file: - shared.env environment: - - EVENTSTORE_GOSSIP_SEED=172.30.240.11:2113,172.30.240.13:2113 + - EVENTSTORE_GOSSIP_SEED=172.30.240.11:1112,172.30.240.13:1112 - EVENTSTORE_REPLICATION_IP=172.30.240.12 - EVENTSTORE_CERTIFICATE_FILE=/etc/eventstore/certs/node/node.crt - EVENTSTORE_CERTIFICATE_PRIVATE_KEY_FILE=/etc/eventstore/certs/node/node.key @@ -70,7 +70,7 @@ services: env_file: - shared.env environment: - - EVENTSTORE_GOSSIP_SEED=172.30.240.11:2113,172.30.240.12:2113 + - EVENTSTORE_GOSSIP_SEED=172.30.240.11:1112,172.30.240.12:1112 - EVENTSTORE_REPLICATION_IP=172.30.240.13 - EVENTSTORE_CERTIFICATE_FILE=/etc/eventstore/certs/node/node.crt - EVENTSTORE_CERTIFICATE_PRIVATE_KEY_FILE=/etc/eventstore/certs/node/node.key diff --git a/proto.lock b/proto.lock index 0534f18460..07fd846f2e 100644 --- a/proto.lock +++ b/proto.lock @@ -8124,4 +8124,4 @@ } } ] -} \ No newline at end of file +} diff --git a/src/EventStore.ClusterNode/Components/Services/ClusterStatusService.cs b/src/EventStore.ClusterNode/Components/Services/ClusterStatusService.cs index 5c753987ec..c8362e60cc 100644 --- a/src/EventStore.ClusterNode/Components/Services/ClusterStatusService.cs +++ b/src/EventStore.ClusterNode/Components/Services/ClusterStatusService.cs @@ -211,7 +211,9 @@ private static ClientClusterInfo.ClientMemberInfo FindMemberByInternalEndpoint( string endpoint) { var cleaned = endpoint.Replace("Unspecified/", "", StringComparison.OrdinalIgnoreCase); - return members.FirstOrDefault(x => string.Equals(InternalTcpEndpoint(x), cleaned, StringComparison.OrdinalIgnoreCase)); + return members.FirstOrDefault(x => + string.Equals(ReplicationEndpoint(x), cleaned, StringComparison.OrdinalIgnoreCase) || + string.Equals(InternalTcpEndpoint(x), cleaned, StringComparison.OrdinalIgnoreCase)); } private static Uri BuildLeaderAddress( @@ -224,6 +226,9 @@ private static string InternalTcpEndpoint(ClientClusterInfo.ClientMemberInfo mem member.InternalTcpIp, member.InternalSecureTcpPort != 0 ? member.InternalSecureTcpPort : member.InternalTcpPort); + private static string ReplicationEndpoint(ClientClusterInfo.ClientMemberInfo member) => + Endpoint(member.ClusterEndPointIp, member.ClusterEndPointPort); + private static string HttpEndpoint(ClientClusterInfo.ClientMemberInfo member) => Endpoint(member.HttpEndPointIp, member.HttpEndPointPort); diff --git a/src/EventStore.ClusterNode/Components/Services/EndpointPolicy.cs b/src/EventStore.ClusterNode/Components/Services/EndpointPolicy.cs new file mode 100644 index 0000000000..bc4382e7fa --- /dev/null +++ b/src/EventStore.ClusterNode/Components/Services/EndpointPolicy.cs @@ -0,0 +1,87 @@ +using System; +using System.Collections.Generic; +using System.Net; +using Google.Protobuf.Reflection; +using Microsoft.AspNetCore.Http; +using Microsoft.AspNetCore.Server.Kestrel.Core; + +namespace EventStore.ClusterNode.Components.Services; + +public enum EndpointRole +{ + Client, + Cluster, +} + +public sealed record EndpointBinding( + EndpointRole Role, + IPEndPoint ListenEndPoint, + HttpProtocols Protocols); + +public sealed record GrpcEndpointRoute(ServiceDescriptor Service, EndpointRole Role) +{ + public PathString Path => new($"/{Service.FullName}"); +} + +public sealed class EndpointPolicy +{ + private readonly IReadOnlyList _bindings; + private readonly IReadOnlyList _routes; + private readonly EndpointRole _defaultRouteRole; + private readonly EndpointRole _nonIpEndpointRole; + + public EndpointPolicy( + IReadOnlyList bindings, + IReadOnlyList routes, + EndpointRole defaultRouteRole, + EndpointRole nonIpEndpointRole) + { + _bindings = bindings; + _routes = routes; + _defaultRouteRole = defaultRouteRole; + _nonIpEndpointRole = nonIpEndpointRole; + } + + public bool Allows(HttpContext context) + { + var endpointRole = GetEndpointRole(context.Connection.LocalIpAddress, context.Connection.LocalPort); + return endpointRole.HasValue && endpointRole.Value == GetRouteRole(context.Request.Path); + } + + private EndpointRole GetRouteRole(PathString requestPath) + { + foreach (var route in _routes) + { + if (requestPath.StartsWithSegments(route.Path, StringComparison.Ordinal)) + return route.Role; + } + + return _defaultRouteRole; + } + + private EndpointRole? GetEndpointRole(IPAddress localAddress, int localPort) + { + if (localAddress is null) + return _nonIpEndpointRole; + + foreach (var binding in _bindings) + { + if (!Matches(binding.ListenEndPoint, localAddress, localPort)) + continue; + + return binding.Role; + } + + return null; + } + + private static bool Matches(IPEndPoint listenEndPoint, IPAddress localAddress, int localPort) + { + if (localPort != listenEndPoint.Port) + return false; + + return listenEndPoint.Address.Equals(IPAddress.Any) || + listenEndPoint.Address.Equals(IPAddress.IPv6Any) || + listenEndPoint.Address.Equals(localAddress); + } +} diff --git a/src/EventStore.ClusterNode/Program.cs b/src/EventStore.ClusterNode/Program.cs index a82abd59cb..e1127f64e2 100644 --- a/src/EventStore.ClusterNode/Program.cs +++ b/src/EventStore.ClusterNode/Program.cs @@ -25,6 +25,7 @@ using Microsoft.AspNetCore.Builder; using Microsoft.AspNetCore.DataProtection; using Microsoft.AspNetCore.Hosting; +using Microsoft.AspNetCore.Http; using Microsoft.AspNetCore.Server.Kestrel.Core; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.DependencyInjection; @@ -272,6 +273,25 @@ async Task Run(ClusterVNodeHostedService hostedService, ManualResetEventSlim sig { x.SuppressStatusMessages = true; }); + EndpointBinding[] endpointBindings = + [ + new(EndpointRole.Client, + new System.Net.IPEndPoint(options.Interface.NodeIp, options.Interface.NodePort), + HttpProtocols.Http1AndHttp2), + new(EndpointRole.Cluster, + options.Interface.GetClusterListenEndPoint(), + HttpProtocols.Http2), + ]; + var endpointPolicy = new EndpointPolicy( + endpointBindings, + [ + new(EventStore.Cluster.Gossip.Descriptor, EndpointRole.Cluster), + new(EventStore.Cluster.Elections.Descriptor, EndpointRole.Cluster), + new(EventStore.Replication.Replication.Descriptor, EndpointRole.Cluster), + new(EventStore.Forwarding.RequestForwarding.Descriptor, EndpointRole.Cluster), + ], + defaultRouteRole: EndpointRole.Client, + nonIpEndpointRole: EndpointRole.Client); builder.WebHost.ConfigureKestrel(server => { @@ -280,9 +300,13 @@ async Task Run(ClusterVNodeHostedService hostedService, ManualResetEventSlim sig server.Limits.Http2.KeepAlivePingTimeout = TimeSpan.FromMilliseconds(options.Grpc.KeepAliveTimeout); - server.Listen(options.Interface.NodeIp, options.Interface.NodePort, listenOptions => - ConfigureHttpOptions(listenOptions, hostedService, - useHttps: !hostedService.Node.DisableHttps)); + foreach (var binding in endpointBindings) + { + server.Listen(binding.ListenEndPoint, listenOptions => + ConfigureHttpOptions(listenOptions, hostedService, + useHttps: !hostedService.Node.DisableHttps, + protocols: binding.Protocols)); + } if (hostedService.Node.EnableUnixSocket) { @@ -325,6 +349,16 @@ async Task Run(ClusterVNodeHostedService hostedService, ManualResetEventSlim sig builder.Services.AddSingleton(hostedService); var app = builder.Build(); + app.Use(async (context, next) => + { + if (!endpointPolicy.Allows(context)) + { + context.Response.StatusCode = StatusCodes.Status404NotFound; + return; + } + + await next(context); + }); app.UseMiddleware(); hostedService.Node.Startup.Configure(app); if (oauthEnabled) @@ -376,14 +410,19 @@ async Task Run(ClusterVNodeHostedService hostedService, ManualResetEventSlim sig } } - private static void ConfigureHttpOptions(ListenOptions listenOptions, ClusterVNodeHostedService hostedService, - bool useHttps) + private static void ConfigureHttpOptions( + ListenOptions listenOptions, + ClusterVNodeHostedService hostedService, + bool useHttps, + HttpProtocols protocols = HttpProtocols.Http1AndHttp2) { + listenOptions.Protocols = protocols; + if (useHttps) { listenOptions.UseHttps(CreateServerOptionsSelectionCallback(hostedService), null); } - else + else if (protocols != HttpProtocols.Http2) { listenOptions.Use(next => new ClearTextHttpMultiplexingMiddleware(next).OnConnectAsync); diff --git a/src/EventStore.Core.Tests/Cluster/MemberInfoTests.cs b/src/EventStore.Core.Tests/Cluster/MemberInfoTests.cs index 0db8ed8764..5852fa8077 100644 --- a/src/EventStore.Core.Tests/Cluster/MemberInfoTests.cs +++ b/src/EventStore.Core.Tests/Cluster/MemberInfoTests.cs @@ -60,12 +60,12 @@ public void member_with_ip_endpoint_should_equal() } [Test] - public void grpc_round_trip_preserves_tcp_and_replication_endpoints() + public void grpc_round_trip_preserves_tcp_and_cluster_endpoints() { var member = EventStore.Core.Cluster.MemberInfo.Initial(Guid.NewGuid(), DateTime.UtcNow, VNodeState.Unknown, true, InternalTcp, null, null, ExternalSecureTcp, Http, - "client", 2113, 1113, 0, false, replicationEndPoint: Replication); + "client", 2113, 1113, 0, false, clusterEndPoint: Replication); var result = FromGrpcClusterInfo(ToGrpcClusterInfo( new EventStore.Core.Cluster.ClusterInfo(member))).Members[0]; @@ -75,11 +75,11 @@ public void grpc_round_trip_preserves_tcp_and_replication_endpoints() Assert.That(result.ExternalTcpEndPoint, Is.Null); Assert.That(result.ExternalSecureTcpEndPoint, Is.EqualTo(ExternalSecureTcp)); Assert.That(result.HttpEndPoint, Is.EqualTo(Http)); - Assert.That(result.ReplicationEndPoint, Is.EqualTo(Replication)); + Assert.That(result.ClusterEndPoint, Is.EqualTo(Replication)); } [Test] - public void explicit_replication_endpoint_is_recognized_without_replacing_tcp_endpoints() + public void explicit_cluster_endpoint_is_recognized_without_replacing_tcp_endpoints() { var member = CreateMember(Replication); var vnode = new VNodeInfo(Guid.NewGuid(), 0, @@ -95,22 +95,42 @@ public void explicit_replication_endpoint_is_recognized_without_replacing_tcp_en Assert.That(member.InternalSecureTcpEndPoint, Is.EqualTo(InternalSecureTcp)); Assert.That(member.ExternalTcpEndPoint, Is.EqualTo(ExternalTcp)); Assert.That(member.ExternalSecureTcpEndPoint, Is.EqualTo(ExternalSecureTcp)); - Assert.That(vnode.ReplicationEndPoint, Is.SameAs(Replication)); - Assert.That(advertise.ReplicationEndPoint, Is.SameAs(Replication)); + Assert.That(vnode.ClusterEndPoint, Is.SameAs(Replication)); + Assert.That(advertise.ClusterEndPoint, Is.SameAs(Replication)); } [Test] - public void client_member_preserves_the_replication_endpoint() + public void client_member_preserves_the_cluster_endpoint() { var clientMember = new EventStore.Core.Cluster.ClientClusterInfo.ClientMemberInfo( CreateMember(Replication)); - Assert.That(clientMember.ReplicationEndPointIp, Is.EqualTo(Replication.Host)); - Assert.That(clientMember.ReplicationEndPointPort, Is.EqualTo(Replication.Port)); + Assert.That(clientMember.ClusterEndPointIp, Is.EqualTo(Replication.Host)); + Assert.That(clientMember.ClusterEndPointPort, Is.EqualTo(Replication.Port)); } [Test] - public void missing_replication_endpoint_falls_back_to_http_endpoint() + public void client_cluster_info_excludes_internal_discovery_placeholders() + { + var member = CreateMember(Replication); + var seed = EventStore.Core.Cluster.MemberInfo.ForManager( + Guid.Empty, + DateTime.UtcNow, + true, + Replication, + clusterEndPoint: Replication); + + var clientCluster = new EventStore.Core.Cluster.ClientClusterInfo( + new EventStore.Core.Cluster.ClusterInfo(member, seed), + Http.Host, + Http.Port); + + Assert.That(clientCluster.Members, Has.Length.EqualTo(1)); + Assert.That(clientCluster.Members[0].InstanceId, Is.EqualTo(member.InstanceId)); + } + + [Test] + public void missing_cluster_endpoint_falls_back_to_http_endpoint() { var member = CreateMember(); var vnode = new VNodeInfo(Guid.NewGuid(), 0, @@ -121,13 +141,13 @@ public void missing_replication_endpoint_falls_back_to_http_endpoint() InternalTcp, InternalSecureTcp, ExternalTcp, ExternalSecureTcp, Http, null, null, 0, null, 0, 0); - Assert.That(member.ReplicationEndPoint, Is.SameAs(Http)); - Assert.That(vnode.ReplicationEndPoint, Is.SameAs(Http)); - Assert.That(advertise.ReplicationEndPoint, Is.SameAs(Http)); + Assert.That(member.ClusterEndPoint, Is.SameAs(Http)); + Assert.That(vnode.ClusterEndPoint, Is.SameAs(Http)); + Assert.That(advertise.ClusterEndPoint, Is.SameAs(Http)); } [Test] - public void grpc_member_without_replication_endpoint_falls_back_to_http_endpoint() + public void grpc_member_without_cluster_endpoint_falls_back_to_http_endpoint() { var grpcCluster = ToGrpcClusterInfo( new EventStore.Core.Cluster.ClusterInfo(CreateMember(Replication))); @@ -135,14 +155,14 @@ public void grpc_member_without_replication_endpoint_falls_back_to_http_endpoint var result = FromGrpcClusterInfo(grpcCluster).Members[0]; - Assert.That(result.ReplicationEndPoint, Is.EqualTo(Http)); + Assert.That(result.ClusterEndPoint, Is.EqualTo(Http)); } - private static EventStore.Core.Cluster.MemberInfo CreateMember(DnsEndPoint replicationEndPoint = null) => + private static EventStore.Core.Cluster.MemberInfo CreateMember(DnsEndPoint clusterEndPoint = null) => EventStore.Core.Cluster.MemberInfo.Initial(Guid.NewGuid(), DateTime.UtcNow, VNodeState.Unknown, true, InternalTcp, InternalSecureTcp, ExternalTcp, ExternalSecureTcp, Http, - "client", 2113, 1113, 0, false, replicationEndPoint: replicationEndPoint); + "client", 2113, 1113, 0, false, clusterEndPoint: clusterEndPoint); private static EventStore.Cluster.ClusterInfo ToGrpcClusterInfo( EventStore.Core.Cluster.ClusterInfo clusterInfo) => diff --git a/src/EventStore.Core.Tests/Helpers/MiniClusterNode.cs b/src/EventStore.Core.Tests/Helpers/MiniClusterNode.cs index 444ab592d2..e84c40b2c8 100644 --- a/src/EventStore.Core.Tests/Helpers/MiniClusterNode.cs +++ b/src/EventStore.Core.Tests/Helpers/MiniClusterNode.cs @@ -25,6 +25,7 @@ using EventStore.Plugins.Subsystems; using EventStore.TcpUnitTestPlugin; using Microsoft.AspNetCore.Hosting; +using Microsoft.AspNetCore.Server.Kestrel.Core; using Microsoft.AspNetCore.Server.Kestrel.Https; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.Hosting; @@ -42,9 +43,9 @@ public class MiniClusterNode private static readonly ILogger Log = Serilog.Log.ForContext>(); - public IPEndPoint InternalTcpEndPoint { get; } public IPEndPoint ExternalTcpEndPoint { get; } public IPEndPoint HttpEndPoint { get; } + public IPEndPoint ClusterEndPoint { get; } public readonly int DebugIndex; @@ -62,22 +63,23 @@ public class MiniClusterNode public VNodeState NodeState = VNodeState.Unknown; private readonly IHost _host; - public MiniClusterNode(string pathname, int debugIndex, IPEndPoint internalTcp, IPEndPoint externalTcp, + public MiniClusterNode(string pathname, int debugIndex, IPEndPoint clusterEndPoint, IPEndPoint externalTcp, IPEndPoint httpEndPoint, EndPoint[] gossipSeeds, ISubsystem[] subsystems = null, bool enableTrustedAuth = false, int memTableSize = 1000, bool disableFlushToDisk = false, bool readOnlyReplica = false, int nodePriority = 0, string intHostAdvertiseAs = null, IExpiryStrategy expiryStrategy = null, ArchiveOptions archiveOptions = null, bool archiver = false, - int clusterSize = 3, bool unsafeAllowSurplusNodes = false) + int clusterSize = 3, bool unsafeAllowSurplusNodes = false, + string replicationHostAdvertiseAs = null) { RunningTime.Start(); RunCount += 1; DebugIndex = debugIndex; - InternalTcpEndPoint = internalTcp; ExternalTcpEndPoint = externalTcp; HttpEndPoint = httpEndPoint; + ClusterEndPoint = clusterEndPoint; _dbPath = Path.Combine( pathname, @@ -117,14 +119,14 @@ public MiniClusterNode(string pathname, int debugIndex, IPEndPoint internalTcp, }, Interface = new() { - ReplicationIp = InternalTcpEndPoint.Address, + ReplicationIp = ClusterEndPoint.Address, NodeIp = ExternalTcpEndPoint.Address, - ReplicationPort = InternalTcpEndPoint.Port, + ReplicationPort = ClusterEndPoint.Port, NodePort = HttpEndPoint.Port, ReplicationHeartbeatTimeout = 2_000, ReplicationHeartbeatInterval = 2_000, EnableTrustedAuth = enableTrustedAuth, - ReplicationHostAdvertiseAs = intHostAdvertiseAs + ReplicationHostAdvertiseAs = replicationHostAdvertiseAs ?? intHostAdvertiseAs }, Database = new() { @@ -215,7 +217,7 @@ public MiniClusterNode(string pathname, int debugIndex, IPEndPoint internalTcp, webHost .UseKestrel(o => { - o.Listen(HttpEndPoint, options => + void ConfigureHttps(ListenOptions options) { options.UseHttps(new HttpsConnectionAdapterOptions { @@ -233,6 +235,13 @@ public MiniClusterNode(string pathname, int debugIndex, IPEndPoint internalTcp, return isValid; } }); + } + + o.Listen(HttpEndPoint, ConfigureHttps); + o.Listen(ClusterEndPoint, options => + { + options.Protocols = HttpProtocols.Http2; + ConfigureHttps(options); }); }) .UseStartup(Node.Startup); diff --git a/src/EventStore.Core.Tests/Integration/Archive/when_archiving_and_restoring_a_cluster.cs b/src/EventStore.Core.Tests/Integration/Archive/when_archiving_and_restoring_a_cluster.cs index 781da7f1ce..888cac4935 100644 --- a/src/EventStore.Core.Tests/Integration/Archive/when_archiving_and_restoring_a_cluster.cs +++ b/src/EventStore.Core.Tests/Integration/Archive/when_archiving_and_restoring_a_cluster.cs @@ -117,7 +117,7 @@ protected override MiniClusterNode CreateNode( new( PathName, index, - endpoints.InternalTcp, + endpoints.ClusterEndPoint, endpoints.ExternalTcp, endpoints.HttpEndPoint, gossipSeeds, @@ -235,7 +235,7 @@ private async Task RestoreNode(int nodeIndex, int coldChunkNumber) private EndPoint[] GossipSeedsFor(int nodeIndex) => _nodeEndpoints .Where((_, index) => index != nodeIndex) - .Select(x => (EndPoint)x.HttpEndPoint) + .Select(x => (EndPoint)x.ClusterEndPoint) .ToArray(); private async Task WaitForArchiveCheckpoint(long minimum) diff --git a/src/EventStore.Core.Tests/Integration/specification_with_cluster.cs b/src/EventStore.Core.Tests/Integration/specification_with_cluster.cs index 6b8aa33884..a5b20c79c8 100644 --- a/src/EventStore.Core.Tests/Integration/specification_with_cluster.cs +++ b/src/EventStore.Core.Tests/Integration/specification_with_cluster.cs @@ -25,13 +25,13 @@ public abstract class specification_with_cluster : Specif protected class Endpoints { - public readonly IPEndPoint InternalTcp; + public readonly IPEndPoint ClusterEndPoint; public readonly IPEndPoint ExternalTcp; public readonly IPEndPoint HttpEndPoint; public IEnumerable Ports() { - yield return InternalTcp.Port; + yield return ClusterEndPoint.Port; yield return ExternalTcp.Port; yield return HttpEndPoint.Port; } @@ -44,9 +44,9 @@ public Endpoints() var defaultLoopBack = new IPEndPoint(IPAddress.Loopback, 0); - var internalTcp = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); - internalTcp.Bind(defaultLoopBack); - _sockets.Add(internalTcp); + var cluster = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); + cluster.Bind(defaultLoopBack); + _sockets.Add(cluster); var externalTcp = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); externalTcp.Bind(defaultLoopBack); @@ -56,7 +56,7 @@ public Endpoints() httpEndPoint.Bind(defaultLoopBack); _sockets.Add(httpEndPoint); - InternalTcp = CopyEndpoint((IPEndPoint)internalTcp.LocalEndPoint); + ClusterEndPoint = CopyEndpoint((IPEndPoint)cluster.LocalEndPoint); ExternalTcp = CopyEndpoint((IPEndPoint)externalTcp.LocalEndPoint); HttpEndPoint = CopyEndpoint((IPEndPoint)httpEndPoint.LocalEndPoint); } @@ -102,7 +102,7 @@ public override async Task TestFixtureSetUp() nodeIndex, _nodeEndpoints[nodeIndex], _nodeEndpoints.Where((_, otherIndex) => otherIndex != nodeIndex) - .Select(x => (EndPoint)x.HttpEndPoint) + .Select(x => (EndPoint)x.ClusterEndPoint) .ToArray(), wait)); _nodes[nodeIndex] = _nodeCreationFactory[nodeIndex](true); @@ -171,7 +171,7 @@ protected virtual void BeforeNodesStart() protected virtual MiniClusterNode CreateNode(int index, Endpoints endpoints, EndPoint[] gossipSeeds, bool wait = true) => new( - PathName, index, endpoints.InternalTcp, + PathName, index, endpoints.ClusterEndPoint, endpoints.ExternalTcp, endpoints.HttpEndPoint, subsystems: Array.Empty(), gossipSeeds: gossipSeeds); diff --git a/src/EventStore.Core.Tests/Integration/when_cluster_nodes_are_restarted.cs b/src/EventStore.Core.Tests/Integration/when_cluster_nodes_are_restarted.cs index 53791c1be1..0d50a0eb6b 100644 --- a/src/EventStore.Core.Tests/Integration/when_cluster_nodes_are_restarted.cs +++ b/src/EventStore.Core.Tests/Integration/when_cluster_nodes_are_restarted.cs @@ -102,7 +102,7 @@ private int SelectRestartNode(bool[] restartedNodes, bool restartLeader) private EndPoint[] GossipSeedsFor(int restartedNodeIndex) => _nodeEndpoints .Where((_, index) => index != restartedNodeIndex) - .Select(x => (EndPoint)x.HttpEndPoint) + .Select(x => (EndPoint)x.ClusterEndPoint) .ToArray(); [Test] diff --git a/src/EventStore.Core.Tests/Integration/when_node_becomes_leader_with_unindexed_data.cs b/src/EventStore.Core.Tests/Integration/when_node_becomes_leader_with_unindexed_data.cs index f7e6400840..2bef053672 100644 --- a/src/EventStore.Core.Tests/Integration/when_node_becomes_leader_with_unindexed_data.cs +++ b/src/EventStore.Core.Tests/Integration/when_node_becomes_leader_with_unindexed_data.cs @@ -25,7 +25,7 @@ namespace EventStore.Core.Tests.Integration; [TestFixture(typeof(LogFormat.V2), typeof(string))] public class when_node_becomes_leader_with_unindexed_data : specification_with_cluster { - private const string FakeHostAdvertiseAs = "192.168.123.123"; + private const string FakeReplicationHostAdvertiseAs = "192.168.123.123"; private const string Username = "admin"; private const string Password = "changeit"; @@ -49,9 +49,9 @@ private class WrongExpectedVersionException : Exception { } protected override async Task Given() { _nodeGossipSeeds = new[] { - new EndPoint[] {_nodeEndpoints[1].HttpEndPoint, _nodeEndpoints[2].HttpEndPoint}, - new EndPoint[] {_nodeEndpoints[0].HttpEndPoint, _nodeEndpoints[2].HttpEndPoint}, - new EndPoint[] {_nodeEndpoints[0].HttpEndPoint, _nodeEndpoints[1].HttpEndPoint} + new EndPoint[] {_nodeEndpoints[1].ClusterEndPoint, _nodeEndpoints[2].ClusterEndPoint}, + new EndPoint[] {_nodeEndpoints[0].ClusterEndPoint, _nodeEndpoints[2].ClusterEndPoint}, + new EndPoint[] {_nodeEndpoints[0].ClusterEndPoint, _nodeEndpoints[1].ClusterEndPoint} }; _httpClient = new HttpClient(new SocketsHttpHandler { @@ -163,15 +163,15 @@ private async Task> ReadAllEvents(IPEndPoint endpoint } private MiniClusterNode CreateNode(int index, Endpoints endpoints, EndPoint[] gossipSeeds, - int nodePriority, string intHostAdvertiseAs) => new( - PathName, index, endpoints.InternalTcp, + int nodePriority, string replicationHostAdvertiseAs) => new( + PathName, index, endpoints.ClusterEndPoint, endpoints.ExternalTcp, endpoints.HttpEndPoint, subsystems: Array.Empty(), gossipSeeds: gossipSeeds, - nodePriority: nodePriority, intHostAdvertiseAs: intHostAdvertiseAs); + nodePriority: nodePriority, replicationHostAdvertiseAs: replicationHostAdvertiseAs); - private Task StartNode(int i, int priority, string intHostAdvertiseAs = null) + private Task StartNode(int i, int priority, string replicationHostAdvertiseAs = null) { - _nodes[i] = CreateNode(i, _nodeEndpoints[i], _nodeGossipSeeds[i], priority, intHostAdvertiseAs); + _nodes[i] = CreateNode(i, _nodeEndpoints[i], _nodeGossipSeeds[i], priority, replicationHostAdvertiseAs); _nodes[i].Start(); return Task.CompletedTask; } @@ -296,9 +296,9 @@ public async Task new_events_should_have_correct_event_numbers(bool appendInitia await ShutdownAllNodes(keepDb: true); // make node 1 become the leader by setting its priority to 1 - // node 0 can't become a follower since it can't replicate over internal TCP due to the fake --int-host-advertise-as - await StartNode(0, priority: 0, intHostAdvertiseAs: FakeHostAdvertiseAs); - await StartNode(1, priority: 1, intHostAdvertiseAs: FakeHostAdvertiseAs); + // node 0 can't become a follower because the fake replication address prevents gRPC replication + await StartNode(0, priority: 0, replicationHostAdvertiseAs: FakeReplicationHostAdvertiseAs); + await StartNode(1, priority: 1, replicationHostAdvertiseAs: FakeReplicationHostAdvertiseAs); try { @@ -336,7 +336,7 @@ public async Task new_events_should_have_correct_event_numbers(bool appendInitia // shut down both nodes await ShutdownAllNodes(maxIdx: 2, keepDb: true); - // start both nodes again without the fake --int-host-advertise-as so that they can form a cluster + // start both nodes again without the fake replication address so that they can form a cluster await StartNode(0, priority: 0); await StartNode(1, priority: 1); try diff --git a/src/EventStore.Core.Tests/Regression/ClusterStatusServiceTests.cs b/src/EventStore.Core.Tests/Regression/ClusterStatusServiceTests.cs new file mode 100644 index 0000000000..b82dc0fe58 --- /dev/null +++ b/src/EventStore.Core.Tests/Regression/ClusterStatusServiceTests.cs @@ -0,0 +1,27 @@ +using System.Collections.Generic; +using System.Reflection; +using EventStore.ClusterNode.Components.Services; +using EventStore.Core.Cluster; +using NUnit.Framework; + +namespace EventStore.Core.Tests.Regression; + +[TestFixture] +public class ClusterStatusServiceTests +{ + [Test] + public void replication_statistics_match_the_member_cluster_endpoint() + { + var member = new ClientClusterInfo.ClientMemberInfo + { + ClusterEndPointIp = "replica.internal", + ClusterEndPointPort = 1112 + }; + + var result = (ClientClusterInfo.ClientMemberInfo)typeof(ClusterStatusService) + .GetMethod("FindMemberByInternalEndpoint", BindingFlags.NonPublic | BindingFlags.Static)! + .Invoke(null, [new List { member }, "replica.internal:1112"])!; + + Assert.That(result, Is.SameAs(member)); + } +} diff --git a/src/EventStore.Core.Tests/Regression/EndpointPolicyTests.cs b/src/EventStore.Core.Tests/Regression/EndpointPolicyTests.cs new file mode 100644 index 0000000000..b60180cca6 --- /dev/null +++ b/src/EventStore.Core.Tests/Regression/EndpointPolicyTests.cs @@ -0,0 +1,124 @@ +using System.Net; +using EventStore.ClusterNode.Components.Services; +using Microsoft.AspNetCore.Http; +using Microsoft.AspNetCore.Server.Kestrel.Core; +using NUnit.Framework; + +namespace EventStore.Core.Tests.Regression; + +[TestFixture] +public class EndpointPolicyTests +{ + [TestCase(2113, "/streams", true)] + [TestCase(2113, "/event_store.client.gossip.Gossip/Read", true)] + [TestCase(2113, "/event_store.cluster.Gossip/Read", false)] + [TestCase(2113, "/event_store.cluster.Elections/Prepare", false)] + [TestCase(2113, "/event_store.replication.Replication/Replicate", false)] + [TestCase(2113, "/event_store.forwarding.RequestForwarding/Forward", false)] + [TestCase(1112, "/streams", false)] + [TestCase(1112, "/event_store.client.gossip.Gossip/Read", false)] + [TestCase(1112, "/event_store.cluster.Gossip/Read", true)] + [TestCase(1112, "/event_store.cluster.Elections/Prepare", true)] + [TestCase(1112, "/event_store.replication.Replication/Replicate", true)] + [TestCase(1112, "/event_store.forwarding.RequestForwarding/Forward", true)] + public void routes_are_isolated_by_endpoint_role(int localPort, string path, bool expected) + { + var policy = CreatePolicy( + new IPEndPoint(IPAddress.Loopback, 2113), + new IPEndPoint(IPAddress.Loopback, 1112)); + var context = new DefaultHttpContext(); + context.Connection.LocalIpAddress = IPAddress.Loopback; + context.Connection.LocalPort = localPort; + context.Request.Path = path; + + Assert.That(policy.Allows(context), Is.EqualTo(expected)); + } + + [Test] + public void wildcard_cluster_binding_matches_the_resolved_local_address() + { + var policy = CreatePolicy( + new IPEndPoint(IPAddress.Loopback, 2113), + new IPEndPoint(IPAddress.Any, 1112)); + var context = new DefaultHttpContext(); + context.Connection.LocalIpAddress = IPAddress.Parse("192.0.2.1"); + context.Connection.LocalPort = 1112; + context.Request.Path = "/event_store.cluster.Gossip/Read"; + + Assert.That(policy.Allows(context), Is.True); + } + + [TestCase("/streams", true)] + [TestCase("/event_store.cluster.Gossip/Read", false)] + public void listeners_without_an_ip_endpoint_remain_client_only(string path, bool expected) + { + var policy = CreatePolicy( + new IPEndPoint(IPAddress.Loopback, 2113), + new IPEndPoint(IPAddress.Loopback, 1112)); + var context = new DefaultHttpContext(); + context.Connection.LocalIpAddress = null; + context.Connection.LocalPort = 0; + context.Request.Path = path; + + Assert.That(policy.Allows(context), Is.EqualTo(expected)); + } + + [Test] + public void additional_routes_can_be_assigned_without_changing_the_policy() + { + var policy = new EndpointPolicy( + [ + new(EndpointRole.Client, new IPEndPoint(IPAddress.Loopback, 2113), HttpProtocols.Http1AndHttp2), + new(EndpointRole.Cluster, new IPEndPoint(IPAddress.Loopback, 1112), HttpProtocols.Http2), + ], + [new(EventStore.Client.Monitoring.Monitoring.Descriptor, EndpointRole.Cluster)], + defaultRouteRole: EndpointRole.Client, + nonIpEndpointRole: EndpointRole.Client); + var context = new DefaultHttpContext(); + context.Connection.LocalIpAddress = IPAddress.Loopback; + context.Connection.LocalPort = 1112; + context.Request.Path = "/event_store.client.monitoring.Monitoring/Stats"; + + Assert.That(policy.Allows(context), Is.True); + } + + [Test] + public void unregistered_ip_listeners_are_denied() + { + var policy = CreatePolicy( + new IPEndPoint(IPAddress.Loopback, 2113), + new IPEndPoint(IPAddress.Loopback, 1112)); + var context = new DefaultHttpContext(); + context.Connection.LocalIpAddress = IPAddress.Loopback; + context.Connection.LocalPort = 3112; + context.Request.Path = "/streams"; + + Assert.That(policy.Allows(context), Is.False); + } + + [Test] + public void bindings_carry_their_transport_protocols() + { + var binding = new EndpointBinding( + EndpointRole.Cluster, + new IPEndPoint(IPAddress.Loopback, 1112), + HttpProtocols.Http2); + + Assert.That(binding.Protocols, Is.EqualTo(HttpProtocols.Http2)); + } + + private static EndpointPolicy CreatePolicy(IPEndPoint clientEndPoint, IPEndPoint clusterEndPoint) => + new( + [ + new(EndpointRole.Client, clientEndPoint, HttpProtocols.Http1AndHttp2), + new(EndpointRole.Cluster, clusterEndPoint, HttpProtocols.Http2), + ], + [ + new(EventStore.Cluster.Gossip.Descriptor, EndpointRole.Cluster), + new(EventStore.Cluster.Elections.Descriptor, EndpointRole.Cluster), + new(EventStore.Replication.Replication.Descriptor, EndpointRole.Cluster), + new(EventStore.Forwarding.RequestForwarding.Descriptor, EndpointRole.Cluster), + ], + defaultRouteRole: EndpointRole.Client, + nonIpEndpointRole: EndpointRole.Client); +} diff --git a/src/EventStore.Core.Tests/Services/ElectionsService/ElectionsServiceTests.cs b/src/EventStore.Core.Tests/Services/ElectionsService/ElectionsServiceTests.cs index bcf9e0938f..16a861db2f 100644 --- a/src/EventStore.Core.Tests/Services/ElectionsService/ElectionsServiceTests.cs +++ b/src/EventStore.Core.Tests/Services/ElectionsService/ElectionsServiceTests.cs @@ -33,7 +33,8 @@ public abstract class ElectionsFixture new IPEndPoint(IPAddress.Loopback, id), new IPEndPoint(IPAddress.Loopback, id), new IPEndPoint(IPAddress.Loopback, id), - new IPEndPoint(IPAddress.Loopback, id), false); + new IPEndPoint(IPAddress.Loopback, id), false, + new IPEndPoint(IPAddress.Loopback, 10_000 + id)); protected static readonly Func MemberInfoFromVNode = (nodeInfo, timestamp, state, isAlive, epochNumber, epochId, priority) => MemberInfo.ForVNode( @@ -42,7 +43,7 @@ public abstract class ElectionsFixture nodeInfo.InternalSecureTcp, nodeInfo.ExternalTcp, nodeInfo.ExternalSecureTcp, nodeInfo.HttpEndPoint, null, 0, 0, 0, 0, 0, 0, epochNumber, epochId, priority, - nodeInfo.IsReadOnlyReplica); + nodeInfo.IsReadOnlyReplica, clusterEndPoint: nodeInfo.ClusterEndPoint); protected ElectionsFixture(VNodeInfo node, VNodeInfo nodeTwo, VNodeInfo nodeThree) { @@ -82,10 +83,10 @@ public void should_send_view_change_to_other_members() _sut.Handle(new ElectionMessage.StartElections()); var expected = new Message[] { - new GrpcMessage.SendOverGrpc(_nodeTwo.HttpEndPoint, + new GrpcMessage.SendOverGrpc(_nodeTwo.ClusterEndPoint, new ElectionMessage.ViewChange(_node.InstanceId, _node.HttpEndPoint, 0), _timeProvider.LocalTime.Add(LeaderElectionProgressTimeout)), - new GrpcMessage.SendOverGrpc(_nodeThree.HttpEndPoint, + new GrpcMessage.SendOverGrpc(_nodeThree.ClusterEndPoint, new ElectionMessage.ViewChange(_node.InstanceId, _node.HttpEndPoint, 0), _timeProvider.LocalTime.Add(LeaderElectionProgressTimeout)), TimerMessage.Schedule.Create( @@ -228,10 +229,10 @@ public void should_send_new_view_change_to_other_members() _sut.Handle(new ElectionMessage.ElectionsTimedOut(view)); var expected = new Message[] { - new GrpcMessage.SendOverGrpc(_nodeTwo.HttpEndPoint, + new GrpcMessage.SendOverGrpc(_nodeTwo.ClusterEndPoint, new ElectionMessage.ViewChange(_node.InstanceId, _node.HttpEndPoint, newView), _timeProvider.LocalTime.Add(LeaderElectionProgressTimeout)), - new GrpcMessage.SendOverGrpc(_nodeThree.HttpEndPoint, + new GrpcMessage.SendOverGrpc(_nodeThree.ClusterEndPoint, new ElectionMessage.ViewChange(_node.InstanceId, _node.HttpEndPoint, newView), _timeProvider.LocalTime.Add(LeaderElectionProgressTimeout)), TimerMessage.Schedule.Create( @@ -320,10 +321,10 @@ public void should_send_view_change_proof_to_other_members() _sut.Handle(new ElectionMessage.SendViewChangeProof()); var expected = new Message[] { - new GrpcMessage.SendOverGrpc(_nodeTwo.HttpEndPoint, + new GrpcMessage.SendOverGrpc(_nodeTwo.ClusterEndPoint, new ElectionMessage.ViewChangeProof(_node.InstanceId, _node.HttpEndPoint, 0), _timeProvider.LocalTime.Add(LeaderElectionProgressTimeout)), - new GrpcMessage.SendOverGrpc(_nodeThree.HttpEndPoint, + new GrpcMessage.SendOverGrpc(_nodeThree.ClusterEndPoint, new ElectionMessage.ViewChangeProof(_node.InstanceId, _node.HttpEndPoint, 0), _timeProvider.LocalTime.Add(LeaderElectionProgressTimeout)), TimerMessage.Schedule.Create( @@ -423,10 +424,10 @@ public void should_send_view_change_to_other_members() _sut.Handle(new ElectionMessage.ViewChange(_nodeTwo.InstanceId, _nodeTwo.HttpEndPoint, 10)); var expected = new Message[] { - new GrpcMessage.SendOverGrpc(_nodeTwo.HttpEndPoint, + new GrpcMessage.SendOverGrpc(_nodeTwo.ClusterEndPoint, new ElectionMessage.ViewChange(_node.InstanceId, _node.HttpEndPoint, 10), _timeProvider.LocalTime.Add(LeaderElectionProgressTimeout)), - new GrpcMessage.SendOverGrpc(_nodeThree.HttpEndPoint, + new GrpcMessage.SendOverGrpc(_nodeThree.ClusterEndPoint, new ElectionMessage.ViewChange(_node.InstanceId, _node.HttpEndPoint, 10), _timeProvider.LocalTime.Add(LeaderElectionProgressTimeout)), TimerMessage.Schedule.Create( @@ -482,10 +483,10 @@ public void should_send_prepares_to_other_members() _sut.Handle(new ElectionMessage.ViewChange(_nodeTwo.InstanceId, _nodeTwo.HttpEndPoint, 0)); var expected = new Message[] { - new GrpcMessage.SendOverGrpc(_nodeTwo.HttpEndPoint, + new GrpcMessage.SendOverGrpc(_nodeTwo.ClusterEndPoint, new ElectionMessage.Prepare(_node.InstanceId, _node.HttpEndPoint, 0), _timeProvider.LocalTime.Add(LeaderElectionProgressTimeout)), - new GrpcMessage.SendOverGrpc(_nodeThree.HttpEndPoint, + new GrpcMessage.SendOverGrpc(_nodeThree.ClusterEndPoint, new ElectionMessage.Prepare(_node.InstanceId, _node.HttpEndPoint, 0), _timeProvider.LocalTime.Add(LeaderElectionProgressTimeout)) }; @@ -602,7 +603,7 @@ public void should_reply_with_prepare_ok() _sut.Handle(new ElectionMessage.Prepare(_nodeTwo.InstanceId, _nodeTwo.HttpEndPoint, 0)); var expected = new Message[] { - new GrpcMessage.SendOverGrpc(_nodeTwo.HttpEndPoint, + new GrpcMessage.SendOverGrpc(_nodeTwo.ClusterEndPoint, new ElectionMessage.PrepareOk(0, _node.InstanceId, _node.HttpEndPoint, -1, -1, Guid.Empty, Guid.Empty, 0, 0, 0, 0, _clusterInfo), _timeProvider.LocalTime.Add(LeaderElectionProgressTimeout)), @@ -732,12 +733,12 @@ public void should_send_proposal_to_other_members() var proposalMessage = (ElectionMessage.Proposal)proposalHttpMessage.Message; var expected = new Message[] { - new GrpcMessage.SendOverGrpc(_nodeTwo.HttpEndPoint, + new GrpcMessage.SendOverGrpc(_nodeTwo.ClusterEndPoint, new ElectionMessage.Proposal(_node.InstanceId, _node.HttpEndPoint, proposalMessage.LeaderId, proposalMessage.LeaderHttpEndPoint, 0, 0, 0, _epochId, Guid.Empty, 0, 0, 0, 0), _timeProvider.LocalTime.Add(LeaderElectionProgressTimeout)), - new GrpcMessage.SendOverGrpc(_nodeThree.HttpEndPoint, + new GrpcMessage.SendOverGrpc(_nodeThree.ClusterEndPoint, new ElectionMessage.Proposal(_node.InstanceId, _node.HttpEndPoint, proposalMessage.LeaderId, proposalMessage.LeaderHttpEndPoint, 0, 0, 0, _epochId, Guid.Empty, 0, 0, 0, 0), @@ -928,12 +929,12 @@ public void should_send_an_acceptance_to_other_members() _nodeThree.InternalSecureTcp, _nodeThree.ExternalTcp, _nodeThree.ExternalSecureTcp, _nodeThree.HttpEndPoint, null, 0, 0, 0, 0, 0, 0, 0, _epochId, 0, _nodeThree.IsReadOnlyReplica)), - new GrpcMessage.SendOverGrpc(_nodeThree.HttpEndPoint, + new GrpcMessage.SendOverGrpc(_nodeThree.ClusterEndPoint, new ElectionMessage.Accept(_node.InstanceId, _node.HttpEndPoint, _nodeThree.InstanceId, _nodeThree.HttpEndPoint, 0), _timeProvider.LocalTime.Add(LeaderElectionProgressTimeout)), - new GrpcMessage.SendOverGrpc(_nodeTwo.HttpEndPoint, + new GrpcMessage.SendOverGrpc(_nodeTwo.ClusterEndPoint, new ElectionMessage.Accept(_node.InstanceId, _node.HttpEndPoint, _nodeThree.InstanceId, _nodeThree.HttpEndPoint, 0), @@ -1195,10 +1196,10 @@ public void should_initiate_leader_resignation_and_inform_other_nodes() _sut.Handle(new ClientMessage.ResignNode()); var expected = new Message[] { - new GrpcMessage.SendOverGrpc(_nodeTwo.HttpEndPoint, + new GrpcMessage.SendOverGrpc(_nodeTwo.ClusterEndPoint, new ElectionMessage.LeaderIsResigning(_node.InstanceId, _node.HttpEndPoint), _timeProvider.LocalTime.Add(LeaderElectionProgressTimeout)), - new GrpcMessage.SendOverGrpc(_nodeThree.HttpEndPoint, + new GrpcMessage.SendOverGrpc(_nodeThree.ClusterEndPoint, new ElectionMessage.LeaderIsResigning(_node.InstanceId, _node.HttpEndPoint), _timeProvider.LocalTime.Add(LeaderElectionProgressTimeout)), }; @@ -1223,7 +1224,41 @@ public void should_reply_with_leader_is_resigning_ok() _sut.Handle(new ElectionMessage.LeaderIsResigning(_nodeTwo.InstanceId, _nodeTwo.HttpEndPoint)); var expected = new Message[] { - new GrpcMessage.SendOverGrpc(_nodeTwo.HttpEndPoint, + new GrpcMessage.SendOverGrpc(_nodeTwo.ClusterEndPoint, + new ElectionMessage.LeaderIsResigningOk( + _nodeTwo.InstanceId, _nodeTwo.HttpEndPoint, + _node.InstanceId, _node.HttpEndPoint), + _timeProvider.LocalTime.Add(LeaderElectionProgressTimeout)), + }; + _publisher.Messages.Should().BeEquivalentTo(expected); + } +} + +public class when_receiving_leader_is_resigning_before_gossip_contains_the_leader : ElectionsFixture +{ + public when_receiving_leader_is_resigning_before_gossip_contains_the_leader() : + base(NodeFactory(1), NodeFactory(2), NodeFactory(3)) + { + _sut.Handle(new GossipMessage.GossipUpdated(new ClusterInfo( + MemberInfoFromVNode(_node, _timeProvider.UtcNow, VNodeState.Unknown, true, 0, _epochId, 0), + MemberInfoFromVNode(_nodeThree, _timeProvider.UtcNow, VNodeState.Unknown, true, 0, _epochId, 0)))); + } + + [Test] + public void should_defer_the_reply_until_gossip_contains_the_leader() + { + _sut.Handle(new ElectionMessage.LeaderIsResigning( + _nodeTwo.InstanceId, + _nodeTwo.HttpEndPoint)); + _publisher.Messages.Should().BeEmpty(); + + _sut.Handle(new GossipMessage.GossipUpdated(new ClusterInfo( + MemberInfoFromVNode(_node, _timeProvider.UtcNow, VNodeState.Unknown, true, 0, _epochId, 0), + MemberInfoFromVNode(_nodeTwo, _timeProvider.UtcNow, VNodeState.Unknown, true, 0, _epochId, 0), + MemberInfoFromVNode(_nodeThree, _timeProvider.UtcNow, VNodeState.Unknown, true, 0, _epochId, 0)))); + + var expected = new Message[] { + new GrpcMessage.SendOverGrpc(_nodeTwo.ClusterEndPoint, new ElectionMessage.LeaderIsResigningOk( _nodeTwo.InstanceId, _nodeTwo.HttpEndPoint, _node.InstanceId, _node.HttpEndPoint), diff --git a/src/EventStore.Core.Tests/Services/GossipService/NodeGossipServiceTests.cs b/src/EventStore.Core.Tests/Services/GossipService/NodeGossipServiceTests.cs index d7c2b10a55..7f230e2f06 100644 --- a/src/EventStore.Core.Tests/Services/GossipService/NodeGossipServiceTests.cs +++ b/src/EventStore.Core.Tests/Services/GossipService/NodeGossipServiceTests.cs @@ -47,32 +47,36 @@ public NodeGossipServiceTestFixture() new IPEndPoint(IPAddress.Loopback, 1111), new IPEndPoint(IPAddress.Loopback, 1111), new IPEndPoint(IPAddress.Loopback, 1111), - new IPEndPoint(IPAddress.Loopback, 1111), false); + new IPEndPoint(IPAddress.Loopback, 1111), false, + new IPEndPoint(IPAddress.Loopback, 11112)); _nodeTwo = new VNodeInfo( Guid.Parse("00000000-0000-0000-0000-000000000002"), 2, new IPEndPoint(IPAddress.Loopback, 2222), new IPEndPoint(IPAddress.Loopback, 2222), new IPEndPoint(IPAddress.Loopback, 2222), new IPEndPoint(IPAddress.Loopback, 2222), - new IPEndPoint(IPAddress.Loopback, 2222), false); + new IPEndPoint(IPAddress.Loopback, 2222), false, + new IPEndPoint(IPAddress.Loopback, 22212)); _nodeThree = new VNodeInfo( Guid.Parse("00000000-0000-0000-0000-000000000003"), 3, new IPEndPoint(IPAddress.Loopback, 3333), new IPEndPoint(IPAddress.Loopback, 3333), new IPEndPoint(IPAddress.Loopback, 3333), new IPEndPoint(IPAddress.Loopback, 3333), - new IPEndPoint(IPAddress.Loopback, 3333), false); + new IPEndPoint(IPAddress.Loopback, 3333), false, + new IPEndPoint(IPAddress.Loopback, 33312)); _nodeFour = new VNodeInfo( Guid.Parse("00000000-0000-0000-0000-000000000004"), 4, new IPEndPoint(IPAddress.Loopback, 4444), new IPEndPoint(IPAddress.Loopback, 4444), new IPEndPoint(IPAddress.Loopback, 4444), new IPEndPoint(IPAddress.Loopback, 4444), - new IPEndPoint(IPAddress.Loopback, 4444), false); + new IPEndPoint(IPAddress.Loopback, 4444), false, + new IPEndPoint(IPAddress.Loopback, 44412)); - _getNodeToGossipTo = infos => infos.First(x => Equals(x.HttpEndPoint, _nodeTwo.HttpEndPoint)); + _getNodeToGossipTo = infos => infos.First(x => Equals(x.ClusterEndPoint, _nodeTwo.ClusterEndPoint)); _gossipSeedSource = new KnownEndpointGossipSeedSource(new[] - {_currentNode.HttpEndPoint, _nodeTwo.HttpEndPoint, _nodeThree.HttpEndPoint}); + {_currentNode.ClusterEndPoint, _nodeTwo.ClusterEndPoint, _nodeThree.ClusterEndPoint}); } [SetUp] @@ -111,7 +115,7 @@ protected Message[] GivenSystemInitializedWithKnownGossipSeedSources(params Mess return new Message[] { new SystemMessage.SystemInit(), new GossipMessage.GotGossipSeedSources(new[] - {_currentNode.HttpEndPoint, _nodeTwo.HttpEndPoint, _nodeThree.HttpEndPoint}) + {_currentNode.ClusterEndPoint, _nodeTwo.ClusterEndPoint, _nodeThree.ClusterEndPoint}) }.Concat(additionalGivens).ToArray(); } @@ -123,7 +127,7 @@ protected static MemberInfo MemberInfoForVNode(VNodeInfo nodeInfo, DateTime utcN nodeInfo.InternalTcp, nodeInfo.InternalSecureTcp, nodeInfo.ExternalTcp, nodeInfo.ExternalSecureTcp, nodeInfo.HttpEndPoint, null, 0, 0, 0, writerCheckpoint ?? 0, 0, -1, epochNumber ?? -1, Guid.Empty, nodePriority ?? 0, false, esVersion, - nodeInfo.ReplicationEndPoint); + nodeInfo.ClusterEndPoint); } /// @@ -131,7 +135,8 @@ protected static MemberInfo MemberInfoForVNode(VNodeInfo nodeInfo, DateTime utcN /// protected static MemberInfo InitialStateForVNode(VNodeInfo nodeInfo, DateTime utcNow, bool isAlive = true, string version = VersionInfo.UnknownVersion) { - return MemberInfo.ForManager(Guid.Empty, utcNow, isAlive, nodeInfo.HttpEndPoint, esVersion: version); + return MemberInfo.ForManager(Guid.Empty, utcNow, isAlive, nodeInfo.ClusterEndPoint, esVersion: version, + clusterEndPoint: nodeInfo.ClusterEndPoint); } } @@ -148,7 +153,7 @@ public void should_get_gossip_sources() { ExpectMessages( new GossipMessage.GotGossipSeedSources(new[] - {_currentNode.HttpEndPoint, _nodeTwo.HttpEndPoint, _nodeThree.HttpEndPoint})); + {_currentNode.ClusterEndPoint, _nodeTwo.ClusterEndPoint, _nodeThree.ClusterEndPoint})); } } @@ -182,7 +187,7 @@ public void should_get_gossip_seeds() { ExpectMessages( new GossipMessage.GotGossipSeedSources(new[] - {_currentNode.HttpEndPoint, _nodeTwo.HttpEndPoint, _nodeThree.HttpEndPoint})); + {_currentNode.ClusterEndPoint, _nodeTwo.ClusterEndPoint, _nodeThree.ClusterEndPoint})); } } @@ -230,25 +235,25 @@ public class when_got_gossip_seed_sources : NodeGossipServiceTestFixture protected override Message When() => new GossipMessage.GotGossipSeedSources(new[] - {_currentNode.HttpEndPoint, _nodeTwo.HttpEndPoint, _nodeThree.HttpEndPoint}); + {_currentNode.ClusterEndPoint, _nodeTwo.ClusterEndPoint, _nodeThree.ClusterEndPoint}); [Test] public void should_start_gossiping_and_schedule_another_gossip() { ExpectMessages( - new GrpcMessage.SendOverGrpc(_nodeTwo.HttpEndPoint, new GossipMessage.SendGossip(new ClusterInfo( + new GrpcMessage.SendOverGrpc(_nodeTwo.ClusterEndPoint, new GossipMessage.SendGossip(new ClusterInfo( MemberInfoForVNode(_currentNode, _timeProvider.UtcNow), InitialStateForVNode(_nodeTwo, _timeProvider.UtcNow), InitialStateForVNode(_nodeThree, _timeProvider.UtcNow)), - _currentNode.HttpEndPoint), _timeProvider.LocalTime.Add(_gossipTimeout)), + _currentNode.ClusterEndPoint), _timeProvider.LocalTime.Add(_gossipTimeout)), TimerMessage.Schedule.Create(GossipServiceBase.GossipStartupInterval, _bus, new GossipMessage.Gossip(1))); } } -public class when_got_gossip_seed_sources_with_distinct_replication_endpoint : NodeGossipServiceTestFixture +public class when_got_gossip_seed_sources_with_distinct_cluster_endpoint : NodeGossipServiceTestFixture { - public when_got_gossip_seed_sources_with_distinct_replication_endpoint() + public when_got_gossip_seed_sources_with_distinct_cluster_endpoint() { _currentNode = new VNodeInfo( Guid.Parse("00000000-0000-0000-0000-000000000001"), 1, @@ -264,13 +269,13 @@ public when_got_gossip_seed_sources_with_distinct_replication_endpoint() protected override Message When() => new GossipMessage.GotGossipSeedSources([ - _currentNode.HttpEndPoint, - _nodeTwo.HttpEndPoint, - _nodeThree.HttpEndPoint + _currentNode.ClusterEndPoint, + _nodeTwo.ClusterEndPoint, + _nodeThree.ClusterEndPoint ]); [Test] - public void should_preserve_the_replication_endpoint() + public void should_preserve_the_cluster_endpoint() { var gossip = (GossipMessage.SendGossip)_bus.Messages .OfType() @@ -278,7 +283,19 @@ public void should_preserve_the_replication_endpoint() .Message; var currentMember = gossip.ClusterInfo.Members.Single(x => x.InstanceId == _currentNode.InstanceId); - Assert.That(currentMember.ReplicationEndPoint, Is.EqualTo(_currentNode.ReplicationEndPoint)); + Assert.That(currentMember.ClusterEndPoint, Is.EqualTo(_currentNode.ClusterEndPoint)); + } + + [Test] + public void should_not_retain_the_cluster_endpoint_seed_as_a_member() + { + var gossip = (GossipMessage.SendGossip)_bus.Messages + .OfType() + .Single() + .Message; + + Assert.That(gossip.ClusterInfo.Members, Has.None.Matches(member => + member.InstanceId == Guid.Empty && member.Is(_currentNode.ClusterEndPoint))); } } @@ -296,11 +313,11 @@ protected override Message When() => public void should_send_the_gossip_over_http_and_schedule_the_next_gossip() { ExpectMessages( - new GrpcMessage.SendOverGrpc(_nodeTwo.HttpEndPoint, new GossipMessage.SendGossip(new ClusterInfo( + new GrpcMessage.SendOverGrpc(_nodeTwo.ClusterEndPoint, new GossipMessage.SendGossip(new ClusterInfo( MemberInfoForVNode(_currentNode, _timeProvider.UtcNow), InitialStateForVNode(_nodeTwo, _timeProvider.UtcNow), InitialStateForVNode(_nodeThree, _timeProvider.UtcNow)), - _currentNode.HttpEndPoint), _timeProvider.LocalTime.Add(_gossipTimeout)), + _currentNode.ClusterEndPoint), _timeProvider.LocalTime.Add(_gossipTimeout)), TimerMessage.Schedule.Create(_gossipInterval, _bus, new GossipMessage.Gossip(++_gossipRound))); } @@ -357,11 +374,11 @@ protected override Message When() => public void should_use_startup_gossip_interval() { ExpectMessages( - new GrpcMessage.SendOverGrpc(_nodeTwo.HttpEndPoint, new GossipMessage.SendGossip(new ClusterInfo( + new GrpcMessage.SendOverGrpc(_nodeTwo.ClusterEndPoint, new GossipMessage.SendGossip(new ClusterInfo( MemberInfoForVNode(_currentNode, _timeProvider.UtcNow), InitialStateForVNode(_nodeTwo, _timeProvider.UtcNow), InitialStateForVNode(_nodeThree, _timeProvider.UtcNow)), - _currentNode.HttpEndPoint), _timeProvider.LocalTime.Add(_gossipTimeout)), + _currentNode.ClusterEndPoint), _timeProvider.LocalTime.Add(_gossipTimeout)), TimerMessage.Schedule.Create(GossipServiceBase.GossipStartupInterval, _bus, new GossipMessage.Gossip(++_gossipRound))); } @@ -381,11 +398,11 @@ protected override Message When() => public void should_use_provided_gossip_interval_for_next_gossip() { ExpectMessages( - new GrpcMessage.SendOverGrpc(_nodeTwo.HttpEndPoint, new GossipMessage.SendGossip(new ClusterInfo( + new GrpcMessage.SendOverGrpc(_nodeTwo.ClusterEndPoint, new GossipMessage.SendGossip(new ClusterInfo( MemberInfoForVNode(_currentNode, _timeProvider.UtcNow), InitialStateForVNode(_nodeTwo, _timeProvider.UtcNow), InitialStateForVNode(_nodeThree, _timeProvider.UtcNow)), - _currentNode.HttpEndPoint), _timeProvider.LocalTime.Add(_gossipTimeout)), + _currentNode.ClusterEndPoint), _timeProvider.LocalTime.Add(_gossipTimeout)), TimerMessage.Schedule.Create(_gossipInterval, _bus, new GossipMessage.Gossip(++_gossipRound))); } @@ -448,7 +465,7 @@ public void gossip_update_must_have_es_version() //updated cluster info should have version info of currentNode, nodeTwo and nodeThree ExpectMessages(new GossipMessage.GossipUpdated(GetExpectedClusterInfo())); //gossip reply should have version info of currentNode, nodeTwo and nodeThree - _capturedMessage.Should().BeEquivalentTo(new GossipMessage.SendGossip(GetExpectedClusterInfo(), _currentNode.HttpEndPoint)); + _capturedMessage.Should().BeEquivalentTo(new GossipMessage.SendGossip(GetExpectedClusterInfo(), _currentNode.ClusterEndPoint)); } private void CaptureGossipReply(Message message) => _capturedMessage = message; @@ -480,7 +497,7 @@ private ClusterInfo GetExpectedClusterInfo() public void reply_should_have_version_info() { _capturedMessage.Should() - .BeEquivalentTo(new GossipMessage.SendGossip(GetExpectedClusterInfo(), _currentNode.HttpEndPoint)); + .BeEquivalentTo(new GossipMessage.SendGossip(GetExpectedClusterInfo(), _currentNode.ClusterEndPoint)); } private void CaptureGossipReply(Message message) => _capturedMessage = message; @@ -731,7 +748,7 @@ protected override Message[] Given() => GivenSystemInitializedWithKnownGossipSeedSources(); protected override Message When() => - new GossipMessage.GossipSendFailed("failed", _nodeTwo.HttpEndPoint); + new GossipMessage.GossipSendFailed("failed", _nodeTwo.ClusterEndPoint); [Test] public void should_mark_the_node_as_dead() @@ -803,7 +820,7 @@ protected override Message When() => public void should_issue_get_gossip() { ExpectMessages( - new GrpcMessage.SendOverGrpc(_currentNode.HttpEndPoint, new GossipMessage.GetGossip(), + new GrpcMessage.SendOverGrpc(_currentNode.ClusterEndPoint, new GossipMessage.GetGossip(), _timeProvider.LocalTime.Add(_gossipTimeout))); } } @@ -1152,6 +1169,36 @@ private static MemberInfo TestNodeFor(int identifier, bool isAlive, DateTime tim private static object[] AllowedNodeRemovalStates => DeadNodeRemoval.AllowedNodeRemovalStates; private static object[] DisallowedNodeRemovalStates => DeadNodeRemoval.DisallowedNodeRemovalStates; + [Test] + public void should_never_replace_self_with_a_newer_cluster_endpoint_seed() + { + var now = DateTime.UtcNow; + var httpEndPoint = new IPEndPoint(IPAddress.Loopback, 2113); + var clusterEndPoint = new IPEndPoint(IPAddress.Loopback, 1112); + var me = MemberInfo.ForVNode( + Guid.NewGuid(), now, VNodeState.Initializing, true, + clusterEndPoint, null, httpEndPoint, null, httpEndPoint, + null, 0, 0, -1, -1, -1, -1, -1, Guid.Empty, 0, false, + clusterEndPoint: clusterEndPoint); + var selfSeed = MemberInfo.ForManager( + Guid.Empty, now.AddSeconds(1), true, clusterEndPoint, + clusterEndPoint: clusterEndPoint); + + var updatedCluster = GossipServiceBase.MergeClusters( + new ClusterInfo(me), + new ClusterInfo(selfSeed), + peerEndPoint: null, + info => info, + now, + me, + currentLeaderInstanceId: null, + allowedTimeDifference: TimeSpan.FromSeconds(1), + deadMemberRemovalTimeout: TimeSpan.FromMinutes(30)); + + Assert.That(updatedCluster.Members, Has.Exactly(1).Matches(member => + member.InstanceId == me.InstanceId && member.ClusterEndPoint.Equals(clusterEndPoint))); + } + [Test] [TestCaseSource(nameof(AllowedNodeRemovalStates))] public void diff --git a/src/EventStore.Core.Tests/Services/Replication/ReadOnlyReplica/connecting_to_read_only_replica.cs b/src/EventStore.Core.Tests/Services/Replication/ReadOnlyReplica/connecting_to_read_only_replica.cs index f8e76a11ab..e59e68f072 100644 --- a/src/EventStore.Core.Tests/Services/Replication/ReadOnlyReplica/connecting_to_read_only_replica.cs +++ b/src/EventStore.Core.Tests/Services/Replication/ReadOnlyReplica/connecting_to_read_only_replica.cs @@ -18,7 +18,7 @@ protected override MiniClusterNode CreateNode(int index, { var isReadOnly = index == 2; var node = new MiniClusterNode( - PathName, index, endpoints.InternalTcp, + PathName, index, endpoints.ClusterEndPoint, endpoints.ExternalTcp, endpoints.HttpEndPoint, gossipSeeds, readOnlyReplica: isReadOnly); if (wait && !isReadOnly) diff --git a/src/EventStore.Core.Tests/Services/RequestForwarding/GrpcRequestForwardingServiceTests.cs b/src/EventStore.Core.Tests/Services/RequestForwarding/GrpcRequestForwardingServiceTests.cs index e861b74159..6534f123e2 100644 --- a/src/EventStore.Core.Tests/Services/RequestForwarding/GrpcRequestForwardingServiceTests.cs +++ b/src/EventStore.Core.Tests/Services/RequestForwarding/GrpcRequestForwardingServiceTests.cs @@ -571,7 +571,7 @@ public void pre_replica_connects_to_the_leader_http_endpoint() fixture.Supervisor.Handle(new SystemMessage.BecomePreReplica( Guid.NewGuid(), Guid.NewGuid(), fixture.Leader)); - Assert.That(fixture.Factory.EndPoints.Single(), Is.EqualTo(fixture.Leader.HttpEndPoint)); + Assert.That(fixture.Factory.EndPoints.Single(), Is.EqualTo(fixture.Leader.ClusterEndPoint)); } [Test] @@ -660,7 +660,7 @@ public void reconnect_keeps_a_healthy_stream_to_the_same_leader() } [Test] - public void reconnect_replaces_the_stream_when_the_http_endpoint_changes() + public void reconnect_keeps_the_stream_when_only_the_client_endpoint_changes() { var fixture = CreateFixture(); fixture.Supervisor.Handle(new SystemMessage.BecomePreReplica( @@ -670,6 +670,24 @@ public void reconnect_replaces_the_stream_when_the_http_endpoint_changes() fixture.Supervisor.Handle(new ReplicationMessage.ReconnectToLeader(Guid.NewGuid(), movedLeader)); + Assert.Multiple(() => + { + Assert.That(fixture.Factory.Services, Has.Exactly(1).Items); + Assert.That(service.StopCalls, Is.Zero); + }); + } + + [Test] + public void reconnect_replaces_the_stream_when_the_cluster_endpoint_changes() + { + var fixture = CreateFixture(); + fixture.Supervisor.Handle(new SystemMessage.BecomePreReplica( + Guid.NewGuid(), Guid.NewGuid(), fixture.Leader)); + var service = fixture.Factory.Services.Single(); + var movedLeader = CreateLeader(fixture.Leader.InstanceId, clusterPort: 3212); + + fixture.Supervisor.Handle(new ReplicationMessage.ReconnectToLeader(Guid.NewGuid(), movedLeader)); + Assert.Multiple(() => { Assert.That(fixture.Factory.Services, Has.Count.EqualTo(2)); @@ -786,7 +804,10 @@ private static Fixture CreateFixture() }; } - private static MemberInfo CreateLeader(Guid? instanceId = null, int httpPort = 2113) => MemberInfo.ForVNode( + private static MemberInfo CreateLeader( + Guid? instanceId = null, + int httpPort = 2113, + int clusterPort = 3112) => MemberInfo.ForVNode( instanceId ?? Guid.NewGuid(), DateTime.UtcNow, VNodeState.Leader, @@ -806,7 +827,8 @@ private static MemberInfo CreateLeader(Guid? instanceId = null, int httpPort = 2 0, Guid.NewGuid(), 0, - false); + false, + clusterEndPoint: new DnsEndPoint("leader-forwarding.internal", clusterPort)); private static async Task WaitUntil(Func condition) { diff --git a/src/EventStore.Core.Tests/Services/RequestForwarding/GrpcRequestForwardingTransportSecurityTests.cs b/src/EventStore.Core.Tests/Services/RequestForwarding/GrpcRequestForwardingTransportSecurityTests.cs index 17fe60ab3f..9851290496 100644 --- a/src/EventStore.Core.Tests/Services/RequestForwarding/GrpcRequestForwardingTransportSecurityTests.cs +++ b/src/EventStore.Core.Tests/Services/RequestForwarding/GrpcRequestForwardingTransportSecurityTests.cs @@ -134,7 +134,8 @@ private static MemberInfo CreateLeader() => MemberInfo.ForVNode( 0, Guid.NewGuid(), 0, - false); + false, + clusterEndPoint: new DnsEndPoint("leader-forwarding.internal", 3112)); private sealed class RejectingServiceFactory : IGrpcRequestForwardingServiceFactory { diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/Replication/GrpcReplicaServiceFactoryTests.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/Replication/GrpcReplicaServiceFactoryTests.cs index 19697515ac..91ba8c4a35 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Grpc/Replication/GrpcReplicaServiceFactoryTests.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/Replication/GrpcReplicaServiceFactoryTests.cs @@ -37,7 +37,30 @@ public void replication_client_uses_node_certificate_names_and_owns_its_http_cli Assert.That(nodeHttpClientFactory.AdditionalCertificateNames, Is.EqualTo(new[] { "cluster.internal" })); client.Dispose(); - Assert.That(nodeHttpClientFactory.Handler.Disposed, Is.True); + Assert.That(nodeHttpClientFactory.TransportHandler.Disposed, Is.True); + } + + [Test] + public void replication_client_normalizes_legacy_sub_second_heartbeat_values() + { + var nodeHttpClientFactory = new RecordingNodeHttpClientFactory(); + var factory = new ReplicationGrpcClientFactory( + Uri.UriSchemeHttps, + nodeHttpClientFactory, + TimeSpan.FromMilliseconds(700), + TimeSpan.FromMilliseconds(700)); + + using var client = factory.Create(new DnsEndPoint("leader.internal", 1112)); + + Assert.Multiple(() => + { + Assert.That(nodeHttpClientFactory.SocketsHandler.KeepAlivePingDelay, + Is.EqualTo(TimeSpan.FromSeconds(1))); + Assert.That(nodeHttpClientFactory.SocketsHandler.KeepAlivePingTimeout, + Is.EqualTo(TimeSpan.FromSeconds(1))); + Assert.That(nodeHttpClientFactory.SocketsHandler.KeepAlivePingPolicy, + Is.EqualTo(HttpKeepAlivePingPolicy.Always)); + }); } [Test] @@ -99,7 +122,8 @@ private sealed class TrackingReplicationGrpcClientFactory : IReplicationGrpcClie private sealed class RecordingNodeHttpClientFactory : INodeHttpClientFactory { - public RecordingHttpMessageHandler Handler { get; } = new(); + public RecordingHttpMessageHandler TransportHandler { get; } = new(); + public SocketsHttpHandler SocketsHandler { get; } = new(); public string[] AdditionalCertificateNames { get; private set; } public HttpClient CreateHttpClient( @@ -107,7 +131,8 @@ public HttpClient CreateHttpClient( Action configureSocketsHttpHandler = null) { AdditionalCertificateNames = additionalCertificateNames; - return new HttpClient(Handler); + configureSocketsHttpHandler?.Invoke(SocketsHandler); + return new HttpClient(TransportHandler); } } diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/Replication/GrpcReplicaServiceSupervisorTests.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/Replication/GrpcReplicaServiceSupervisorTests.cs index a6eaedf281..8c3c96bb6f 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Grpc/Replication/GrpcReplicaServiceSupervisorTests.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/Replication/GrpcReplicaServiceSupervisorTests.cs @@ -22,7 +22,7 @@ public class GrpcReplicaServiceSupervisorTests [TestCase(false)] [TestCase(true)] - public async Task pre_replica_state_starts_and_tracks_a_stream_for_the_advertised_http_endpoints( + public async Task pre_replica_state_starts_and_tracks_a_stream_for_the_advertised_cluster_endpoints( bool readOnlyReplica) { var fixture = CreateFixture(); @@ -37,7 +37,7 @@ public async Task pre_replica_state_starts_and_tracks_a_stream_for_the_advertise var request = fixture.Factory.Requests.Single(); Assert.Multiple(() => { - Assert.That(request.Endpoints.LeaderEndPoint, Is.EqualTo(fixture.Leader.HttpEndPoint)); + Assert.That(request.Endpoints.LeaderEndPoint, Is.EqualTo(fixture.Leader.ClusterEndPoint)); Assert.That(request.Endpoints.AdvertisedReplicaEndPoint, Is.EqualTo(fixture.AdvertisedEndPoint)); Assert.That(request.Service.StartCalls, Is.EqualTo(1)); Assert.That(fixture.TrackedTasks.Single(), Is.SameAs(request.Service.Task)); @@ -410,7 +410,7 @@ private static Fixture CreateFixture( startException, createException, beforeCreateReturns); - var advertisedEndPoint = new DnsEndPoint("replica.internal", 2113); + var advertisedEndPoint = new DnsEndPoint("replica.internal", 1112); var trackedTasks = new List(); var supervisor = new GrpcReplicaServiceSupervisor( publisher, @@ -447,7 +447,8 @@ private static MemberInfo CreateLeader() => MemberInfo.ForVNode( 0, Guid.NewGuid(), 0, - false); + false, + clusterEndPoint: new DnsEndPoint("leader.replication.internal", 1112)); private static SystemMessage.StateChangeMessage CreateReplicaState( VNodeState state, diff --git a/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_building/with_cluster_node_and_custom_settings.cs b/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_building/with_cluster_node_and_custom_settings.cs index e9eb8027df..769e24f24a 100644 --- a/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_building/with_cluster_node_and_custom_settings.cs +++ b/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_building/with_cluster_node_and_custom_settings.cs @@ -153,8 +153,8 @@ public void should_advertise_the_configured_replication_port() { Assert.Multiple(() => { - Assert.AreEqual(_options.Interface.ReplicationPort, _node.NodeInfo.ReplicationEndPoint.GetPort()); - Assert.AreEqual(3112, _node.GossipAdvertiseInfo.ReplicationEndPoint.Port); + Assert.AreEqual(_options.Interface.ReplicationPort, _node.NodeInfo.ClusterEndPoint.GetPort()); + Assert.AreEqual(3112, _node.GossipAdvertiseInfo.ClusterEndPoint.Port); Assert.AreEqual(3112, _node.GossipAdvertiseInfo.InternalSecureTcp.Port); }); } diff --git a/src/EventStore.Core.XUnit.Tests/Configuration/ClusterVNodeOptionsTests.cs b/src/EventStore.Core.XUnit.Tests/Configuration/ClusterVNodeOptionsTests.cs index f18a6018cf..4425678ae6 100644 --- a/src/EventStore.Core.XUnit.Tests/Configuration/ClusterVNodeOptionsTests.cs +++ b/src/EventStore.Core.XUnit.Tests/Configuration/ClusterVNodeOptionsTests.cs @@ -102,6 +102,7 @@ public void replication_port_advertise_as_is_configurable() options.Interface.ReplicationPortAdvertiseAs.Should().Be(2112); options.Interface.GetReplicationPortAdvertiseAs().Should().Be(2112); + options.Interface.GetClusterPortAdvertiseAs().Should().Be(2112); Assert.Empty(options.Unknown.Options); } @@ -119,6 +120,7 @@ public void replication_port_advertise_as_is_configurable_from_the_environment() options.Interface.ReplicationPortAdvertiseAs.Should().Be(2112); options.Interface.GetReplicationPortAdvertiseAs().Should().Be(2112); + options.Interface.GetClusterPortAdvertiseAs().Should().Be(2112); } [Fact] @@ -137,6 +139,7 @@ public void replication_port_advertise_as_is_configurable_from_yaml() options.Interface.ReplicationPortAdvertiseAs.Should().Be(2112); options.Interface.GetReplicationPortAdvertiseAs().Should().Be(2112); + options.Interface.GetClusterPortAdvertiseAs().Should().Be(2112); } finally { @@ -150,6 +153,7 @@ public void deprecated_replication_tcp_port_advertise_as_remains_compatible() var options = GetOptions("--replication-tcp-port-advertise-as 3112"); options.Interface.GetReplicationPortAdvertiseAs().Should().Be(3112); + options.Interface.GetClusterPortAdvertiseAs().Should().Be(3112); options.GetDeprecationWarnings().Should().Contain( "ReplicationTcpPortAdvertiseAs setting has been deprecated"); Assert.Empty(options.Unknown.Options); @@ -162,6 +166,7 @@ public void replication_port_advertise_as_takes_precedence_over_deprecated_alias "--replication-port-advertise-as 2112 --replication-tcp-port-advertise-as 3112"); options.Interface.GetReplicationPortAdvertiseAs().Should().Be(2112); + options.Interface.GetClusterPortAdvertiseAs().Should().Be(2112); } [Fact] diff --git a/src/EventStore.Core.XUnit.Tests/Configuration/ClusterVNodeOptionsValidatorTests.cs b/src/EventStore.Core.XUnit.Tests/Configuration/ClusterVNodeOptionsValidatorTests.cs index 3caee89cc6..37a3703c9c 100644 --- a/src/EventStore.Core.XUnit.Tests/Configuration/ClusterVNodeOptionsValidatorTests.cs +++ b/src/EventStore.Core.XUnit.Tests/Configuration/ClusterVNodeOptionsValidatorTests.cs @@ -1,4 +1,5 @@ using System; +using System.Net; using EventStore.Common.Exceptions; using EventStore.Core.Authentication; using EventStore.Core.Services; @@ -9,6 +10,69 @@ namespace EventStore.Core.XUnit.Tests.Configuration; // Some other tests are in ClusterNodeOptionsTests/when_building public class ClusterVNodeOptionsValidatorTests { + [Theory] + [InlineData("127.0.0.1", "127.0.0.1")] + [InlineData("0.0.0.0", "127.0.0.1")] + [InlineData("127.0.0.1", "0.0.0.0")] + public void node_and_replication_listeners_cannot_overlap(string nodeIp, string replicationIp) + { + var options = new ClusterVNodeOptions + { + Interface = new() + { + NodeIp = IPAddress.Parse(nodeIp), + NodePort = 2113, + ReplicationIp = IPAddress.Parse(replicationIp), + ReplicationPort = 2113, + } + }; + + Assert.Throws(() => ClusterVNodeOptionsValidator.Validate(options)); + } + + [Fact] + public void node_and_replication_listeners_can_use_the_same_port_on_distinct_addresses() + { + var options = new ClusterVNodeOptions + { + Interface = new() + { + NodeIp = IPAddress.Parse("127.0.0.1"), + NodePort = 2113, + ReplicationIp = IPAddress.Parse("127.0.0.2"), + ReplicationPort = 2113, + } + }; + + ClusterVNodeOptionsValidator.Validate(options); + } + + [Theory] + [InlineData(0)] + [InlineData(-1)] + public void replication_heartbeat_interval_must_be_positive(int interval) + { + var options = new ClusterVNodeOptions + { + Interface = new() { ReplicationHeartbeatInterval = interval } + }; + + Assert.Throws(() => ClusterVNodeOptionsValidator.Validate(options)); + } + + [Theory] + [InlineData(0)] + [InlineData(-1)] + public void replication_heartbeat_timeout_must_be_positive(int timeout) + { + var options = new ClusterVNodeOptions + { + Interface = new() { ReplicationHeartbeatTimeout = timeout } + }; + + Assert.Throws(() => ClusterVNodeOptionsValidator.Validate(options)); + } + [Theory] [InlineData(false, false, true)] [InlineData(false, true, true)] diff --git a/src/EventStore.Core/Cluster/ClientClusterInfo.cs b/src/EventStore.Core/Cluster/ClientClusterInfo.cs index b1e516d8f8..573b762c42 100644 --- a/src/EventStore.Core/Cluster/ClientClusterInfo.cs +++ b/src/EventStore.Core/Cluster/ClientClusterInfo.cs @@ -17,7 +17,10 @@ public ClientClusterInfo() public ClientClusterInfo(ClusterInfo clusterInfo, string serverIp, int serverPort) { - Members = clusterInfo.Members.Select(x => new ClientMemberInfo(x)).ToArray(); + Members = clusterInfo.Members + .Where(x => x.State != VNodeState.Manager) + .Select(x => new ClientMemberInfo(x)) + .ToArray(); ServerIp = serverIp; ServerPort = serverPort; } @@ -47,8 +50,8 @@ public class ClientMemberInfo public string InternalHttpEndPointIp { get; set; } public int InternalHttpEndPointPort { get; set; } - public string ReplicationEndPointIp { get; set; } - public int ReplicationEndPointPort { get; set; } + public string ClusterEndPointIp { get; set; } + public int ClusterEndPointPort { get; set; } public string HttpEndPointIp { get; set; } public int HttpEndPointPort { get; set; } @@ -86,8 +89,8 @@ public ClientMemberInfo(MemberInfo member) InternalHttpEndPointIp = member.HttpEndPoint.GetHost(); InternalHttpEndPointPort = member.HttpEndPoint.GetPort(); - ReplicationEndPointIp = member.ReplicationEndPoint.GetHost(); - ReplicationEndPointPort = member.ReplicationEndPoint.GetPort(); + ClusterEndPointIp = member.ClusterEndPoint.GetHost(); + ClusterEndPointPort = member.ClusterEndPoint.GetPort(); HttpEndPointIp = string.IsNullOrEmpty(member.AdvertiseHostToClientAs) ? member.HttpEndPoint.GetHost() @@ -129,7 +132,7 @@ public override string ToString() $"InternalTcpIp: {InternalTcpIp}, InternalTcpPort: {InternalTcpPort}, InternalSecureTcpPort: {InternalSecureTcpPort}, " + $"ExternalTcpIp: {ExternalTcpIp}, ExternalTcpPort: {ExternalTcpPort}, ExternalSecureTcpPort: {ExternalSecureTcpPort}, " + $"InternalHttpEndPointIp: {InternalHttpEndPointIp}, InternalHttpEndPointPort: {InternalHttpEndPointPort}, " + - $"ReplicationEndPointIp: {ReplicationEndPointIp}, ReplicationEndPointPort: {ReplicationEndPointPort}, " + + $"ClusterEndPointIp: {ClusterEndPointIp}, ClusterEndPointPort: {ClusterEndPointPort}, " + $"HttpEndPointIp: {HttpEndPointIp}, HttpEndPointPort: {HttpEndPointPort}, " + $"LastCommitPosition: {LastCommitPosition}, WriterCheckpoint: {WriterCheckpoint}, ChaserCheckpoint: {ChaserCheckpoint}, " + $"EpochPosition: {EpochPosition}, EpochNumber: {EpochNumber}, EpochId: {EpochId:B}, NodePriority: {NodePriority}, " + diff --git a/src/EventStore.Core/Cluster/ClusterInfo.cs b/src/EventStore.Core/Cluster/ClusterInfo.cs index 6276119ed3..48870b2e29 100644 --- a/src/EventStore.Core/Cluster/ClusterInfo.cs +++ b/src/EventStore.Core/Cluster/ClusterInfo.cs @@ -104,8 +104,8 @@ internal static EventStore.Cluster.ClusterInfo ToGrpcClusterInfo(ClusterInfo clu x.HttpEndPoint.GetHost(), (uint)x.HttpEndPoint.GetPort()), ReplicationEndPoint = new EventStore.Cluster.EndPoint( - x.ReplicationEndPoint.GetHost(), - (uint)x.ReplicationEndPoint.GetPort()), + x.ClusterEndPoint.GetHost(), + (uint)x.ClusterEndPoint.GetPort()), InternalTcp = x.InternalSecureTcpEndPoint != null ? new EventStore.Cluster.EndPoint( x.InternalSecureTcpEndPoint.GetHost(), diff --git a/src/EventStore.Core/Cluster/MemberInfo.cs b/src/EventStore.Core/Cluster/MemberInfo.cs index fbcacf901c..248ac6ee5d 100644 --- a/src/EventStore.Core/Cluster/MemberInfo.cs +++ b/src/EventStore.Core/Cluster/MemberInfo.cs @@ -20,7 +20,8 @@ public class MemberInfo : IEquatable public readonly EndPoint ExternalTcpEndPoint; public readonly EndPoint ExternalSecureTcpEndPoint; public readonly EndPoint HttpEndPoint; - public readonly EndPoint ReplicationEndPoint; + public readonly EndPoint ClusterEndPoint; + public EndPoint ReplicationEndPoint => ClusterEndPoint; public readonly string AdvertiseHostToClientAs; public readonly int AdvertiseHttpPortToClientAs; public readonly int AdvertiseTcpPortToClientAs; @@ -39,12 +40,12 @@ public class MemberInfo : IEquatable public static MemberInfo ForManager(Guid instanceId, DateTime timeStamp, bool isAlive, EndPoint httpEndPoint, string esVersion = VersionInfo.UnknownVersion, - EndPoint replicationEndPoint = null) + EndPoint clusterEndPoint = null) { return new MemberInfo(instanceId, timeStamp, VNodeState.Manager, isAlive, httpEndPoint, null, httpEndPoint, null, httpEndPoint, null, 0, 0, - -1, -1, -1, -1, -1, Guid.Empty, 0, false, esVersion, replicationEndPoint); + -1, -1, -1, -1, -1, Guid.Empty, 0, false, esVersion, clusterEndPoint); } public static MemberInfo ForVNode(Guid instanceId, @@ -67,7 +68,7 @@ public static MemberInfo ForVNode(Guid instanceId, Guid epochId, int nodePriority, bool isReadOnlyReplica, string esVersion = VersionInfo.UnknownVersion, - EndPoint replicationEndPoint = null) + EndPoint clusterEndPoint = null) { if (state == VNodeState.Manager) { @@ -80,7 +81,7 @@ public static MemberInfo ForVNode(Guid instanceId, httpEndPoint, advertiseHostToClientAs, advertiseHttpPortToClientAs, advertiseTcpPortToClientAs, lastCommitPosition, writerCheckpoint, chaserCheckpoint, epochPosition, epochNumber, epochId, nodePriority, isReadOnlyReplica, esVersion, - replicationEndPoint); + clusterEndPoint); } public static MemberInfo Initial(Guid instanceId, @@ -97,7 +98,7 @@ public static MemberInfo Initial(Guid instanceId, int advertiseTcpPortToClientAs, int nodePriority, bool isReadOnlyReplica, string esVersion = VersionInfo.UnknownVersion, - EndPoint replicationEndPoint = null) + EndPoint clusterEndPoint = null) { if (state == VNodeState.Manager) { @@ -109,7 +110,7 @@ public static MemberInfo Initial(Guid instanceId, externalTcpEndPoint, externalSecureTcpEndPoint, httpEndPoint, advertiseHostToClientAs, advertiseHttpPortToClientAs, advertiseTcpPortToClientAs, -1, -1, -1, -1, -1, Guid.Empty, nodePriority, isReadOnlyReplica, esVersion, - replicationEndPoint); + clusterEndPoint); } internal MemberInfo(Guid instanceId, DateTime timeStamp, VNodeState state, bool isAlive, @@ -118,7 +119,7 @@ internal MemberInfo(Guid instanceId, DateTime timeStamp, VNodeState state, bool EndPoint httpEndPoint, string advertiseHostToClientAs, int advertiseHttpPortToClientAs, int advertiseTcpPortToClientAs, long lastCommitPosition, long writerCheckpoint, long chaserCheckpoint, long epochPosition, int epochNumber, Guid epochId, int nodePriority, bool isReadOnlyReplica, - string esVersion = null, EndPoint replicationEndPoint = null) + string esVersion = null, EndPoint clusterEndPoint = null) { Ensure.Equal(false, internalTcpEndPoint == null && internalSecureTcpEndPoint == null, "Both internal TCP endpoints are null"); Ensure.NotNull(httpEndPoint, nameof(httpEndPoint)); @@ -134,7 +135,7 @@ internal MemberInfo(Guid instanceId, DateTime timeStamp, VNodeState state, bool ExternalTcpEndPoint = externalTcpEndPoint; ExternalSecureTcpEndPoint = externalSecureTcpEndPoint; HttpEndPoint = httpEndPoint; - ReplicationEndPoint = replicationEndPoint ?? httpEndPoint; + ClusterEndPoint = clusterEndPoint ?? httpEndPoint; AdvertiseHostToClientAs = advertiseHostToClientAs; AdvertiseHttpPortToClientAs = advertiseHttpPortToClientAs; AdvertiseTcpPortToClientAs = advertiseTcpPortToClientAs; @@ -168,7 +169,7 @@ internal MemberInfo(MemberInfoDto dto) ? new DnsEndPoint(dto.ExternalTcpIp, dto.ExternalSecureTcpPort) : null; HttpEndPoint = new DnsEndPoint(dto.HttpEndPointIp, dto.HttpEndPointPort); - ReplicationEndPoint = HttpEndPoint; + ClusterEndPoint = HttpEndPoint; AdvertiseHostToClientAs = dto.AdvertiseHostToClientAs; AdvertiseHttpPortToClientAs = dto.AdvertiseHttpPortToClientAs; AdvertiseTcpPortToClientAs = dto.AdvertiseTcpPortToClientAs; @@ -186,7 +187,7 @@ public bool Is(EndPoint endPoint) { return endPoint != null && (HttpEndPoint.EndPointEquals(endPoint) - || ReplicationEndPoint.EndPointEquals(endPoint) + || ClusterEndPoint.EndPointEquals(endPoint) || (InternalTcpEndPoint != null && InternalTcpEndPoint.EndPointEquals(endPoint)) || (InternalSecureTcpEndPoint != null && InternalSecureTcpEndPoint.EndPointEquals(endPoint)) || (ExternalTcpEndPoint != null && ExternalTcpEndPoint.EndPointEquals(endPoint)) @@ -221,7 +222,7 @@ public MemberInfo Updated(DateTime utcNow, epoch != null ? epoch.EpochNumber : EpochNumber, epoch != null ? epoch.EpochId : EpochId, nodePriority ?? NodePriority, - IsReadOnlyReplica, esVersion ?? ESVersion, ReplicationEndPoint); + IsReadOnlyReplica, esVersion ?? ESVersion, ClusterEndPoint); } public override string ToString() @@ -238,7 +239,7 @@ public override string ToString() $"{(InternalSecureTcpEndPoint == null ? "n/a" : InternalSecureTcpEndPoint.ToString())}, " + $"{(ExternalTcpEndPoint == null ? "n/a" : ExternalTcpEndPoint.ToString())}, " + $"{(ExternalSecureTcpEndPoint == null ? "n/a" : ExternalSecureTcpEndPoint.ToString())}, " + - $"Replication:{ReplicationEndPoint}, " + + $"Cluster:{ClusterEndPoint}, " + $"{HttpEndPoint}, (ADVERTISED: HTTP:{AdvertiseHostToClientAs}:{AdvertiseHttpPortToClientAs}, TCP:{AdvertiseHostToClientAs}:{AdvertiseTcpPortToClientAs}), " + $"Version: {ESVersion}] " + $"{LastCommitPosition}/{WriterCheckpoint}/{ChaserCheckpoint}/E{EpochNumber}@{EpochPosition}:{EpochId:B} | {TimeStamp:yyyy-MM-dd HH:mm:ss.fff}"; @@ -265,7 +266,7 @@ public bool Equals(MemberInfo other) && Equals(other.ExternalTcpEndPoint, ExternalTcpEndPoint) && Equals(other.ExternalSecureTcpEndPoint, ExternalSecureTcpEndPoint) && Equals(other.HttpEndPoint, HttpEndPoint) - && Equals(other.ReplicationEndPoint, ReplicationEndPoint) + && Equals(other.ClusterEndPoint, ClusterEndPoint) && other.AdvertiseHostToClientAs == AdvertiseHostToClientAs && other.AdvertiseHttpPortToClientAs == AdvertiseHttpPortToClientAs && other.AdvertiseTcpPortToClientAs == AdvertiseTcpPortToClientAs @@ -311,7 +312,7 @@ public override int GetHashCode() result = (result * 397) ^ (ExternalSecureTcpEndPoint != null ? ExternalSecureTcpEndPoint.GetHashCode() : 0); result = (result * 397) ^ HttpEndPoint.GetHashCode(); - result = (result * 397) ^ ReplicationEndPoint.GetHashCode(); + result = (result * 397) ^ ClusterEndPoint.GetHashCode(); result = (result * 397) ^ (AdvertiseHostToClientAs != null ? AdvertiseHostToClientAs.GetHashCode() : 0); result = (result * 397) ^ AdvertiseHttpPortToClientAs.GetHashCode(); result = (result * 397) ^ AdvertiseTcpPortToClientAs.GetHashCode(); diff --git a/src/EventStore.Core/ClusterVNode.cs b/src/EventStore.Core/ClusterVNode.cs index 0ac48b9d39..041c45e00f 100644 --- a/src/EventStore.Core/ClusterVNode.cs +++ b/src/EventStore.Core/ClusterVNode.cs @@ -297,9 +297,7 @@ public ClusterVNode(ClusterVNodeOptions options, var enableExternalTcp = nodeTcpOptions.EnableExternalTcp; var httpEndPoint = new IPEndPoint(options.Interface.NodeIp, options.Interface.NodePort); - var replicationEndPoint = new IPEndPoint( - options.Interface.ReplicationIp, - options.Interface.ReplicationPort); + var clusterEndPoint = options.Interface.GetClusterListenEndPoint(); var intTcp = disableInternalTcpTls ? new IPEndPoint(options.Interface.ReplicationIp, @@ -319,9 +317,9 @@ public ClusterVNode(ClusterVNodeOptions options, nodeTcpOptions.NodeTcpPort) : null; - var replicationPortAdvertiseAs = options.Interface.GetReplicationPortAdvertiseAs(); - var intTcpPortAdvertiseAs = disableInternalTcpTls ? replicationPortAdvertiseAs : 0; - var intSecTcpPortAdvertiseAs = !disableInternalTcpTls ? replicationPortAdvertiseAs : 0; + var clusterPortAdvertiseAs = options.Interface.GetClusterPortAdvertiseAs(); + var intTcpPortAdvertiseAs = disableInternalTcpTls ? clusterPortAdvertiseAs : 0; + var intSecTcpPortAdvertiseAs = !disableInternalTcpTls ? clusterPortAdvertiseAs : 0; var extTcpPortAdvertiseAs = enableExternalTcp && disableExternalTcpTls && nodeTcpOptions.NodeTcpPortAdvertiseAs.HasValue @@ -335,7 +333,7 @@ public ClusterVNode(ClusterVNodeOptions options, Log.Information("Quorum size set to {quorum}.", options.Cluster.QuorumSize); NodeInfo = new VNodeInfo(instanceId.Value, debugIndex, intTcp, intSecIp, extTcp, extSecIp, - httpEndPoint, options.Cluster.ReadOnlyReplica, replicationEndPoint); + httpEndPoint, options.Cluster.ReadOnlyReplica, clusterEndPoint); var metricsConfiguration = MetricsConfiguration.Get(configuration); var trackers = new Trackers(); @@ -922,7 +920,8 @@ GossipAdvertiseInfo GetGossipAdvertiseInfo() IPAddress extIpAddress = options.Interface.NodeIp; - var intHostToAdvertise = options.Interface.ReplicationHostAdvertiseAs ?? intIpAddress.ToString(); + var clusterHostAdvertiseAs = options.Interface.GetClusterHostAdvertiseAs(); + var intHostToAdvertise = clusterHostAdvertiseAs ?? intIpAddress.ToString(); var extHostToAdvertise = options.Interface.NodeHostAdvertiseAs ?? extIpAddress.ToString(); if (intIpAddress.Equals(IPAddress.Any) || extIpAddress.Equals(IPAddress.Any)) @@ -931,7 +930,7 @@ GossipAdvertiseInfo GetGossipAdvertiseInfo() IPAddress addressToAdvertise = options.Cluster.ClusterSize > 1 ? nonLoopbackAddress : IPAddress.Loopback; - if (intIpAddress.Equals(IPAddress.Any) && options.Interface.ReplicationHostAdvertiseAs == null) + if (intIpAddress.Equals(IPAddress.Any) && clusterHostAdvertiseAs == null) { intHostToAdvertise = addressToAdvertise.ToString(); } @@ -970,16 +969,16 @@ GossipAdvertiseInfo GetGossipAdvertiseInfo() options.Interface.NodePortAdvertiseAs > 0 ? options.Interface.NodePortAdvertiseAs : NodeInfo.HttpEndPoint.GetPort()); - var advertisedReplicationEndPoint = new DnsEndPoint(intHostToAdvertise, - replicationPortAdvertiseAs > 0 - ? replicationPortAdvertiseAs - : NodeInfo.ReplicationEndPoint.GetPort()); + var advertisedClusterEndPoint = new DnsEndPoint(intHostToAdvertise, + clusterPortAdvertiseAs > 0 + ? clusterPortAdvertiseAs + : NodeInfo.ClusterEndPoint.GetPort()); return new GossipAdvertiseInfo(intTcpEndPoint, intSecureTcpEndPoint, extTcpEndPoint, - extSecureTcpEndPoint, httpEndPoint, options.Interface.ReplicationHostAdvertiseAs, + extSecureTcpEndPoint, httpEndPoint, clusterHostAdvertiseAs, options.Interface.NodeHostAdvertiseAs, options.Interface.NodePortAdvertiseAs, options.Interface.AdvertiseHostToClientAs, options.Interface.AdvertiseNodePortToClientAs, - nodeTcpOptions?.NodeTcpPortAdvertiseAs ?? 0, advertisedReplicationEndPoint); + nodeTcpOptions?.NodeTcpPortAdvertiseAs ?? 0, advertisedClusterEndPoint); } _httpService = new KestrelHttpService(_mainQueue, NodeInfo.HttpEndPoint); @@ -1493,7 +1492,7 @@ GossipAdvertiseInfo GetGossipAdvertiseInfo() GossipAdvertiseInfo.AdvertiseHttpPortToClientAs, GossipAdvertiseInfo.AdvertiseTcpPortToClientAs, options.Cluster.NodePriority, options.Cluster.ReadOnlyReplica, VersionInfo.Version, - GossipAdvertiseInfo.ReplicationEndPoint); + GossipAdvertiseInfo.ClusterEndPoint); // ELECTIONS TRACKER _mainBus.Subscribe(trackers.ElectionCounterTracker); @@ -1537,13 +1536,17 @@ GossipAdvertiseInfo GetGossipAdvertiseInfo() _grpcReplicaServiceSupervisor = new GrpcReplicaServiceSupervisor( _mainQueue, new GrpcReplicaServiceFactory( - new ReplicationGrpcClientFactory(uriScheme, _nodeHttpClientFactory), + new ReplicationGrpcClientFactory( + uriScheme, + _nodeHttpClientFactory, + TimeSpan.FromMilliseconds(options.Interface.ReplicationHeartbeatInterval), + TimeSpan.FromMilliseconds(options.Interface.ReplicationHeartbeatTimeout)), new ReplicaSubscriptionDataSource(Db, epochManager), NodeInfo.InstanceId, options.Cluster.ReadOnlyReplica ? ReplicaPromotability.NonPromotable : ReplicaPromotability.Promotable), - GossipAdvertiseInfo.HttpEndPoint, + GossipAdvertiseInfo.ClusterEndPoint, AddTask); _mainBus.Subscribe(_grpcReplicaServiceSupervisor); _mainBus.Subscribe(_grpcReplicaServiceSupervisor); @@ -1593,13 +1596,18 @@ GossipAdvertiseInfo GetGossipAdvertiseInfo() // GOSSIP + var clusterGossipPort = options.Cluster.ClusterGossipPort > 0 + ? options.Cluster.ClusterGossipPort + : clusterPortAdvertiseAs > 0 + ? clusterPortAdvertiseAs + : options.Interface.ReplicationPort; var gossipSeedSource = ( options.Cluster.DiscoverViaDns, options.Cluster.ClusterSize > 1, options.Cluster.GossipSeed is { Length: > 0 }) switch { (true, true, _) => (IGossipSeedSource)new DnsGossipSeedSource(options.Cluster.ClusterDns, - options.Cluster.ClusterGossipPort), + clusterGossipPort), (false, true, false) => throw new InvalidConfigurationException( "DNS discovery is disabled, but no gossip seed endpoints have been specified. " + "Specify gossip seeds using the `GossipSeed` option."), @@ -2272,5 +2280,6 @@ private void ReloadCertificates(ClusterVNodeOptions options) } public override string ToString() => - $"[{NodeInfo.InstanceId:B}, {NodeInfo.InternalTcp}, {NodeInfo.ExternalTcp}, {NodeInfo.HttpEndPoint}]"; + $"[{NodeInfo.InstanceId:B}, {NodeInfo.InternalTcp}, {NodeInfo.ExternalTcp}, " + + $"{NodeInfo.ClusterEndPoint}, {NodeInfo.HttpEndPoint}]"; } diff --git a/src/EventStore.Core/Configuration/ClusterVNodeOptions.cs b/src/EventStore.Core/Configuration/ClusterVNodeOptions.cs index 620253566d..2dca0e5ad0 100644 --- a/src/EventStore.Core/Configuration/ClusterVNodeOptions.cs +++ b/src/EventStore.Core/Configuration/ClusterVNodeOptions.cs @@ -408,10 +408,10 @@ public record ClusterOptions [Description("DNS name from which other nodes can be discovered.")] public string ClusterDns { get; init; } = "fake.dns"; - [Description("The port on which cluster nodes' managers are running.")] - public int ClusterGossipPort { get; init; } = 2113; + [Description("The internal cluster port used for DNS gossip discovery. A value of 0 uses the advertised cluster port.")] + public int ClusterGossipPort { get; init; } = 0; - [Description("Endpoints for other cluster nodes from which to seed gossip.")] + [Description("Internal cluster endpoints for other nodes from which to seed gossip.")] public EndPoint[] GossipSeed { get; init; } = []; [Description("The interval, in ms, nodes should try to gossip with each other."), @@ -620,7 +620,7 @@ public record GrpcOptions [Description("Interface Options")] public record InterfaceOptions { - [Description("The IP Address used by internal replication between nodes in the cluster.")] + [Description("The IP address used by the internal gRPC cluster listener. The option name is retained for compatibility.")] public IPAddress ReplicationIp { get; init; } = IPAddress.Loopback; [Description("The IP Address for the node.")] @@ -629,13 +629,13 @@ public record InterfaceOptions [Description("The Port to run the HTTP server on.")] public int NodePort { get; init; } = 2113; - [Description("The TCP port used by internal replication between nodes in the cluster.")] + [Description("The port used by the internal gRPC cluster listener. The option name is retained for compatibility.")] public int ReplicationPort { get; init; } = 1112; [Description("Advertise the Node's host name to other nodes and external clients as.")] public string? NodeHostAdvertiseAs { get; init; } = null; - [Description("Advertise the Replication host name to other nodes in the cluster as.")] + [Description("Advertise the internal cluster host name to other nodes. The option name is retained for compatibility.")] public string? ReplicationHostAdvertiseAs { get; init; } = null; [Description("Advertise Host in Gossip to Client As.")] @@ -647,10 +647,10 @@ public record InterfaceOptions [Description("Advertise Http Port As.")] public int NodePortAdvertiseAs { get; init; } = 0; - [Description("Advertise the gRPC replication port as.")] + [Description("Advertise the internal gRPC cluster port. The option name is retained for compatibility.")] public int ReplicationPortAdvertiseAs { get; init; } = 0; - [Description("Deprecated alias for ReplicationPortAdvertiseAs. Advertise the gRPC replication port as.")] + [Description("Deprecated alias for ReplicationPortAdvertiseAs. Advertise the internal gRPC cluster port.")] [Deprecated( "The ReplicationTcpPortAdvertiseAs setting has been deprecated because replication uses gRPC. " + "Use ReplicationPortAdvertiseAs instead.")] @@ -661,11 +661,17 @@ public int GetReplicationPortAdvertiseAs() => ? ReplicationPortAdvertiseAs : ReplicationTcpPortAdvertiseAs; - [Description("Heartbeat timeout for Replication TCP sockets."), + public IPEndPoint GetClusterListenEndPoint() => new(ReplicationIp, ReplicationPort); + + public string? GetClusterHostAdvertiseAs() => ReplicationHostAdvertiseAs; + + public int GetClusterPortAdvertiseAs() => GetReplicationPortAdvertiseAs(); + + [Description("Keepalive ping timeout for internal gRPC cluster connections. Values below 1000 ms use the HTTP/2 minimum of 1000 ms."), Unit("ms")] public int ReplicationHeartbeatTimeout { get; init; } = 700; - [Description("Heartbeat interval for Replication TCP sockets."), + [Description("Keepalive ping interval for internal gRPC cluster connections. Values below 1000 ms use the HTTP/2 minimum of 1000 ms."), Unit("ms")] public int ReplicationHeartbeatInterval { get; init; } = 700; diff --git a/src/EventStore.Core/Configuration/ClusterVNodeOptionsExtensions.cs b/src/EventStore.Core/Configuration/ClusterVNodeOptionsExtensions.cs index 562cbcd82a..052f6bb17a 100644 --- a/src/EventStore.Core/Configuration/ClusterVNodeOptionsExtensions.cs +++ b/src/EventStore.Core/Configuration/ClusterVNodeOptionsExtensions.cs @@ -99,7 +99,7 @@ public static ClusterVNodeOptions WithExternalTcpOn( options with { Interface = options.Interface with { NodeIp = endPoint.Address, } }; /// - /// Sets the internal tcp endpoint to the specified value + /// Sets the gRPC replication endpoint to the specified value /// /// The /// The internal endpoint to use @@ -148,6 +148,18 @@ options with } }; + public static ClusterVNodeOptions AdvertiseReplicationHostAs( + this ClusterVNodeOptions options, + EndPoint endPoint) => + options with + { + Interface = options.Interface with + { + ReplicationHostAdvertiseAs = endPoint.GetHost(), + ReplicationPortAdvertiseAs = endPoint.GetPort() + } + }; + /// /// /// The diff --git a/src/EventStore.Core/Configuration/ClusterVNodeOptionsValidator.cs b/src/EventStore.Core/Configuration/ClusterVNodeOptionsValidator.cs index 41c7b4b39f..35d4a1312a 100644 --- a/src/EventStore.Core/Configuration/ClusterVNodeOptionsValidator.cs +++ b/src/EventStore.Core/Configuration/ClusterVNodeOptionsValidator.cs @@ -32,6 +32,23 @@ public static void Validate(ClusterVNodeOptions options) throw new ArgumentNullException(nameof(options.Interface.ReplicationIp)); } + if (options.Interface.NodePort == options.Interface.ReplicationPort && + EndpointsOverlap(options.Interface.NodeIp, options.Interface.ReplicationIp)) + { + throw new ArgumentException( + $"{nameof(options.Interface.NodePort)} and {nameof(options.Interface.ReplicationPort)} cannot bind the same endpoint."); + } + + if (options.Interface.ReplicationHeartbeatInterval <= 0) + { + throw new ArgumentOutOfRangeException(nameof(options.Interface.ReplicationHeartbeatInterval)); + } + + if (options.Interface.ReplicationHeartbeatTimeout <= 0) + { + throw new ArgumentOutOfRangeException(nameof(options.Interface.ReplicationHeartbeatTimeout)); + } + if (options.Cluster.ClusterSize <= 0) { throw new ArgumentOutOfRangeException(nameof(options.Cluster.ClusterSize), options.Cluster.ClusterSize, @@ -174,4 +191,11 @@ public static bool ValidateForStartup(ClusterVNodeOptions options) return true; } + private static bool EndpointsOverlap(System.Net.IPAddress first, System.Net.IPAddress second) => + first.Equals(second) || + first.Equals(System.Net.IPAddress.Any) || + second.Equals(System.Net.IPAddress.Any) || + first.Equals(System.Net.IPAddress.IPv6Any) || + second.Equals(System.Net.IPAddress.IPv6Any); + } diff --git a/src/EventStore.Core/Data/GossipAdvertiseInfo.cs b/src/EventStore.Core/Data/GossipAdvertiseInfo.cs index 4fcfcca51f..471959a2c9 100644 --- a/src/EventStore.Core/Data/GossipAdvertiseInfo.cs +++ b/src/EventStore.Core/Data/GossipAdvertiseInfo.cs @@ -10,7 +10,8 @@ public class GossipAdvertiseInfo public DnsEndPoint ExternalTcp { get; } public DnsEndPoint ExternalSecureTcp { get; } public DnsEndPoint HttpEndPoint { get; } - public DnsEndPoint ReplicationEndPoint { get; } + public DnsEndPoint ClusterEndPoint { get; } + public DnsEndPoint ReplicationEndPoint => ClusterEndPoint; public string AdvertiseInternalHostAs { get; } public string AdvertiseExternalHostAs { get; } public int AdvertiseHttpPortAs { get; } @@ -23,7 +24,7 @@ public GossipAdvertiseInfo(DnsEndPoint internalTcp, DnsEndPoint internalSecureTc DnsEndPoint httpEndPoint, string advertiseInternalHostAs, string advertiseExternalHostAs, int advertiseHttpPortAs, string advertiseHostToClientAs, int advertiseHttpPortToClientAs, int advertiseTcpPortToClientAs, - DnsEndPoint replicationEndPoint = null) + DnsEndPoint clusterEndPoint = null) { Ensure.Equal(false, internalTcp == null && internalSecureTcp == null, "Both internal TCP endpoints are null"); @@ -32,7 +33,7 @@ public GossipAdvertiseInfo(DnsEndPoint internalTcp, DnsEndPoint internalSecureTc ExternalTcp = externalTcp; ExternalSecureTcp = externalSecureTcp; HttpEndPoint = httpEndPoint; - ReplicationEndPoint = replicationEndPoint ?? httpEndPoint; + ClusterEndPoint = clusterEndPoint ?? httpEndPoint; AdvertiseInternalHostAs = advertiseInternalHostAs; AdvertiseExternalHostAs = advertiseExternalHostAs; AdvertiseHttpPortAs = advertiseHttpPortAs; @@ -46,7 +47,7 @@ public override string ToString() return string.Format( $"IntTcp: {InternalTcp}, IntSecureTcp: {InternalSecureTcp}\n" + $"ExtTcp: {ExternalTcp}, ExtSecureTcp: {ExternalSecureTcp}\n" + - $"Http: {HttpEndPoint}, Replication: {ReplicationEndPoint}, HttpAdvertiseAs: {AdvertiseExternalHostAs}:{AdvertiseHttpPortAs},\n" + + $"Http: {HttpEndPoint}, Cluster: {ClusterEndPoint}, HttpAdvertiseAs: {AdvertiseExternalHostAs}:{AdvertiseHttpPortAs},\n" + $"HttpAdvertiseToClientAs: {AdvertiseHostToClientAs}:{AdvertiseHttpPortToClientAs},\n" + $"TcpAdvertiseToClientAs: {AdvertiseHostToClientAs}:{AdvertiseTcpPortToClientAs}"); } diff --git a/src/EventStore.Core/Data/VNodeInfo.cs b/src/EventStore.Core/Data/VNodeInfo.cs index d8aa86d500..53f328dc7d 100644 --- a/src/EventStore.Core/Data/VNodeInfo.cs +++ b/src/EventStore.Core/Data/VNodeInfo.cs @@ -13,7 +13,8 @@ public class VNodeInfo public readonly IPEndPoint ExternalTcp; public readonly IPEndPoint ExternalSecureTcp; public readonly EndPoint HttpEndPoint; - public readonly EndPoint ReplicationEndPoint; + public readonly EndPoint ClusterEndPoint; + public EndPoint ReplicationEndPoint => ClusterEndPoint; public readonly bool IsReadOnlyReplica; public VNodeInfo(Guid instanceId, int debugIndex, @@ -21,7 +22,7 @@ public VNodeInfo(Guid instanceId, int debugIndex, IPEndPoint externalTcp, IPEndPoint externalSecureTcp, EndPoint httpEndPoint, bool isReadOnlyReplica, - EndPoint replicationEndPoint = null) + EndPoint clusterEndPoint = null) { Ensure.NotEmptyGuid(instanceId, "instanceId"); Ensure.Equal(false, internalTcp == null && internalSecureTcp == null, "Both internal TCP endpoints are null"); @@ -34,7 +35,7 @@ public VNodeInfo(Guid instanceId, int debugIndex, ExternalTcp = externalTcp; ExternalSecureTcp = externalSecureTcp; HttpEndPoint = httpEndPoint; - ReplicationEndPoint = replicationEndPoint ?? httpEndPoint; + ClusterEndPoint = clusterEndPoint ?? httpEndPoint; IsReadOnlyReplica = isReadOnlyReplica; } @@ -42,7 +43,7 @@ public bool Is(EndPoint endPoint) { return endPoint != null && (HttpEndPoint.Equals(endPoint) - || ReplicationEndPoint.Equals(endPoint) + || ClusterEndPoint.Equals(endPoint) || (InternalTcp != null && InternalTcp.Equals(endPoint)) || (InternalSecureTcp != null && InternalSecureTcp.Equals(endPoint)) || (ExternalTcp != null && ExternalTcp.Equals(endPoint)) @@ -53,14 +54,14 @@ public override string ToString() { return string.Format("InstanceId: {0:B}, InternalTcp: {1}, InternalSecureTcp: {2}, " + "ExternalTcp: {3}, ExternalSecureTcp: {4}, HttpEndPoint: {5}, " + - "ReplicationEndPoint: {6}, IsReadOnlyReplica: {7}", + "ClusterEndPoint: {6}, IsReadOnlyReplica: {7}", InstanceId, InternalTcp, InternalSecureTcp, ExternalTcp, ExternalSecureTcp, HttpEndPoint, - ReplicationEndPoint, + ClusterEndPoint, IsReadOnlyReplica); } } diff --git a/src/EventStore.Core/Services/ElectionsService.cs b/src/EventStore.Core/Services/ElectionsService.cs index d865b66e26..3a1b29c1b5 100644 --- a/src/EventStore.Core/Services/ElectionsService.cs +++ b/src/EventStore.Core/Services/ElectionsService.cs @@ -77,8 +77,10 @@ public class ElectionsService : IHandle, private Guid? _leader; private Guid? _lastElectedLeader; - private MemberInfo[] _servers; + private MemberInfo[] _clusterMembership; + private MemberInfo[] _liveClusterMembers; private Guid? _resigningLeaderInstanceId; + private ElectionMessage.LeaderIsResigning _pendingLeaderResignation; public ElectionsService(IPublisher publisher, MemberInfo memberInfo, @@ -132,7 +134,7 @@ public ElectionsService(IPublisher publisher, } var ownInfo = GetOwnInfo(); - _servers = new[] { + _liveClusterMembers = new[] { MemberInfo.ForVNode(memberInfo.InstanceId, _timeProvider.UtcNow, VNodeState.Initializing, @@ -143,8 +145,9 @@ public ElectionsService(IPublisher publisher, memberInfo.AdvertiseHostToClientAs, memberInfo.AdvertiseHttpPortToClientAs, memberInfo.AdvertiseTcpPortToClientAs, ownInfo.LastCommitPosition, ownInfo.WriterCheckpoint, ownInfo.ChaserCheckpoint, ownInfo.EpochPosition, ownInfo.EpochNumber, ownInfo.EpochId, ownInfo.NodePriority, - memberInfo.IsReadOnlyReplica, VersionInfo.Version) + memberInfo.IsReadOnlyReplica, VersionInfo.Version, memberInfo.ClusterEndPoint) }; + _clusterMembership = _liveClusterMembers; } public void SubscribeMessages(ISubscriber subscriber) @@ -207,14 +210,22 @@ public void Handle(ElectionMessage.LeaderIsResigning message) { Log.Information("ELECTIONS: LEADER IS RESIGNING [{leaderHttpEndPoint}, {leaderId:B}].", message.LeaderHttpEndPoint, message.LeaderId); + var leader = _clusterMembership.FirstOrDefault(x => x.InstanceId == message.LeaderId); + if (leader is null) + { + _pendingLeaderResignation = message; + return; + } + _pendingLeaderResignation = null; + var leaderIsResigningMessageOk = new ElectionMessage.LeaderIsResigningOk( message.LeaderId, - message.LeaderHttpEndPoint, + leader.HttpEndPoint, _memberInfo.InstanceId, _memberInfo.HttpEndPoint); _resigningLeaderInstanceId = message.LeaderId; - _publisher.Publish(new GrpcMessage.SendOverGrpc(message.LeaderHttpEndPoint, leaderIsResigningMessageOk, + _publisher.Publish(new GrpcMessage.SendOverGrpc(leader.ClusterEndPoint, leaderIsResigningMessageOk, _timeProvider.LocalTime.Add(_leaderElectionProgressTimeout))); } @@ -243,10 +254,17 @@ public void Handle(SystemMessage.BecomeShuttingDown message) public void Handle(GossipMessage.GossipUpdated message) { - _servers = message.ClusterInfo.Members.Where(x => x.State != VNodeState.Manager) + _clusterMembership = message.ClusterInfo.Members + .Where(x => x.State != VNodeState.Manager) + .ToArray(); + _liveClusterMembers = _clusterMembership .Where(x => x.IsAlive) .OrderByDescending(x => x.HttpEndPoint, IPComparer) .ToArray(); + if (_pendingLeaderResignation is not null) + { + Handle(_pendingLeaderResignation); + } } public void Handle(ElectionMessage.StartElections message) @@ -309,9 +327,9 @@ private void ShiftToLeaderElection(int view) private void SendToAllExceptMe(Message message) { - foreach (var server in _servers.Where(x => x.InstanceId != _memberInfo.InstanceId)) + foreach (var server in _liveClusterMembers.Where(x => x.InstanceId != _memberInfo.InstanceId)) { - _publisher.Publish(new GrpcMessage.SendOverGrpc(server.HttpEndPoint, message, + _publisher.Publish(new GrpcMessage.SendOverGrpc(server.ClusterEndPoint, message, _timeProvider.LocalTime.Add(_leaderElectionProgressTimeout))); } } @@ -412,7 +430,7 @@ public void Handle(ElectionMessage.ViewChangeProof message) private bool AmILeaderOf(int lastAttemptedView) { - var serversExcludingNonPotentialLeaders = _servers.Where(x => !x.IsReadOnlyReplica).ToArray(); + var serversExcludingNonPotentialLeaders = _liveClusterMembers.Where(x => !x.IsReadOnlyReplica).ToArray(); var leader = serversExcludingNonPotentialLeaders[lastAttemptedView % serversExcludingNonPotentialLeaders.Length]; return leader.InstanceId == _memberInfo.InstanceId; @@ -447,7 +465,8 @@ public void Handle(ElectionMessage.Prepare message) return; } - if (_servers.All(x => x.InstanceId != message.ServerId)) + var server = _liveClusterMembers.FirstOrDefault(x => x.InstanceId == message.ServerId); + if (server is null) { return; // unknown instance } @@ -461,14 +480,14 @@ public void Handle(ElectionMessage.Prepare message) } var prepareOk = CreatePrepareOk(message.View); - _publisher.Publish(new GrpcMessage.SendOverGrpc(message.ServerHttpEndPoint, prepareOk, + _publisher.Publish(new GrpcMessage.SendOverGrpc(server.ClusterEndPoint, prepareOk, _timeProvider.LocalTime.Add(_leaderElectionProgressTimeout))); } private ElectionMessage.PrepareOk CreatePrepareOk(int view) { var ownInfo = GetOwnInfo(); - var clusterInfo = new ClusterInfo(_servers); + var clusterInfo = new ClusterInfo(_liveClusterMembers); return new ElectionMessage.PrepareOk(view, ownInfo.InstanceId, ownInfo.HttpEndPoint, ownInfo.EpochNumber, ownInfo.EpochPosition, ownInfo.EpochId, ownInfo.EpochLeaderInstanceId, ownInfo.LastCommitPosition, ownInfo.WriterCheckpoint, ownInfo.ChaserCheckpoint, @@ -528,7 +547,11 @@ private void SendProposal() _acceptsReceived.Clear(); _leaderProposal = null; - var leader = GetBestLeaderCandidate(_prepareOkReceived, _servers, _resigningLeaderInstanceId, _lastAttemptedView); + var leader = GetBestLeaderCandidate( + _prepareOkReceived, + _liveClusterMembers, + _resigningLeaderInstanceId, + _lastAttemptedView); if (leader == null) { Log.Information("ELECTIONS: (V={lastAttemptedView}) NO LEADER CANDIDATE WHEN TRYING TO SEND PROPOSAL.", @@ -746,12 +769,12 @@ public void Handle(ElectionMessage.Proposal message) return; } - if (_servers.All(x => x.InstanceId != message.ServerId)) + if (_liveClusterMembers.All(x => x.InstanceId != message.ServerId)) { return; } - if (_servers.All(x => x.InstanceId != message.LeaderId)) + if (_liveClusterMembers.All(x => x.InstanceId != message.LeaderId)) { return; } @@ -779,7 +802,7 @@ public void Handle(ElectionMessage.Proposal message) var ownInfo = GetOwnInfo(); if (!IsLegitimateLeader(message.View, message.ServerHttpEndPoint, message.ServerId, - candidate, _servers, _lastElectedLeader, _memberInfo.InstanceId, ownInfo, + candidate, _liveClusterMembers, _lastElectedLeader, _memberInfo.InstanceId, ownInfo, _resigningLeaderInstanceId)) { return; @@ -841,7 +864,7 @@ public void Handle(ElectionMessage.Accept message) if (_acceptsReceived.Add(message.ServerId) && _acceptsReceived.Count == _clusterSize / 2 + 1) { - var leader = _servers.FirstOrDefault(x => x.InstanceId == _leaderProposal.InstanceId); + var leader = _liveClusterMembers.FirstOrDefault(x => x.InstanceId == _leaderProposal.InstanceId); if (leader != null) { _leader = _leaderProposal.InstanceId; diff --git a/src/EventStore.Core/Services/Gossip/GossipServiceBase.cs b/src/EventStore.Core/Services/Gossip/GossipServiceBase.cs index 643147e6ed..94e006e256 100644 --- a/src/EventStore.Core/Services/Gossip/GossipServiceBase.cs +++ b/src/EventStore.Core/Services/Gossip/GossipServiceBase.cs @@ -166,8 +166,8 @@ public void Handle(GossipMessage.Gossip message) { _cluster = UpdateCluster(_cluster, x => x.InstanceId == _memberInfo.InstanceId ? GetUpdatedMe(x) : x, _timeProvider, DeadMemberRemovalPeriod, CurrentRole); - _bus.Publish(new GrpcMessage.SendOverGrpc(node.HttpEndPoint, - new GossipMessage.SendGossip(_cluster, _memberInfo.HttpEndPoint), + _bus.Publish(new GrpcMessage.SendOverGrpc(node.ClusterEndPoint, + new GossipMessage.SendGossip(_cluster, _memberInfo.ClusterEndPoint), _timeProvider.LocalTime.Add(GossipTimeout))); } @@ -213,7 +213,7 @@ public void Handle(GossipMessage.GossipReceived message) _timeProvider.UtcNow, _memberInfo, CurrentLeader?.InstanceId, AllowedTimeDifference, DeadMemberRemovalPeriod); - message.Envelope.ReplyWith(new GossipMessage.SendGossip(_cluster, _memberInfo.HttpEndPoint)); + message.Envelope.ReplyWith(new GossipMessage.SendGossip(_cluster, _memberInfo.ClusterEndPoint)); if (_cluster.HasChangedSince(oldCluster)) { @@ -227,7 +227,7 @@ public void Handle(GossipMessage.ReadGossip message) { if (_cluster != null) { - message.Envelope.ReplyWith(new GossipMessage.SendGossip(_cluster, _memberInfo.HttpEndPoint)); + message.Envelope.ReplyWith(new GossipMessage.SendGossip(_cluster, _memberInfo.ClusterEndPoint)); } } @@ -304,7 +304,7 @@ public void Handle(SystemMessage.VNodeConnectionLost message) Log.Information("Looks like node [{nodeEndPoint}] is DEAD (TCP connection lost). Issuing a gossip to confirm.", message.VNodeEndPoint); - _bus.Publish(new GrpcMessage.SendOverGrpc(node.HttpEndPoint, + _bus.Publish(new GrpcMessage.SendOverGrpc(node.ClusterEndPoint, new GossipMessage.GetGossip(), _timeProvider.LocalTime.Add(GossipTimeout))); } @@ -405,12 +405,14 @@ public static ClusterInfo MergeClusters(ClusterInfo myCluster, ClusterInfo other bool isPeerOld = peerNode?.ESVersion == null; foreach (var member in othersCluster.Members) { - if (member.InstanceId == me.InstanceId || member.Is(me.HttpEndPoint) - ) // we know about ourselves better + if (member.InstanceId == me.InstanceId || + member.Is(me.HttpEndPoint) || + member.Is(me.ClusterEndPoint)) // we know about ourselves better { continue; } + var existingMem = members.Values.FirstOrDefault(existing => IsSameMember(existing, member)); if (member.Equals(peerNode)) // peer knows about itself better { if ((utcNow - member.TimeStamp).Duration() > allowedTimeDifference) @@ -419,35 +421,35 @@ public static ClusterInfo MergeClusters(ClusterInfo myCluster, ClusterInfo other + "UTC now: {dateTime:yyyy-MM-dd HH:mm:ss.fff}, peer's time stamp: {peerTimestamp:yyyy-MM-dd HH:mm:ss.fff}.", peerEndPoint, utcNow, member.TimeStamp); } - members[member.HttpEndPoint] = member.Updated(utcNow: member.TimeStamp, esVersion: isPeerOld ? VersionInfo.OldVersion : member.ESVersion); + SetMember(members, existingMem, + member.Updated(utcNow: member.TimeStamp, + esVersion: isPeerOld ? VersionInfo.OldVersion : member.ESVersion)); } else { - MemberInfo existingMem; // if there is no data about this member or data is stale -- update - if (!members.TryGetValue(member.HttpEndPoint, out existingMem) || - IsMoreUpToDate(member, existingMem)) + if (existingMem is null || IsMoreUpToDate(member, existingMem)) { // we do not trust leader's alive status and state to come from outside if (currentLeaderInstanceId != null && existingMem != null && member.InstanceId == currentLeaderInstanceId) { - members[member.HttpEndPoint] = + SetMember(members, existingMem, member.Updated(utcNow: utcNow, isAlive: existingMem.IsAlive, - state: existingMem.State); + state: existingMem.State)); } else { - members[member.HttpEndPoint] = member; + SetMember(members, existingMem, member); } } if (peerNode != null && isPeerOld) { - MemberInfo newInfo = members[member.HttpEndPoint]; + var newInfo = members.Values.First(x => IsSameMember(x, member)); // if we don't have past information about es version of the node, es version is unknown because old peer won't be sending version info in gossip - members[member.HttpEndPoint] = newInfo.Updated(newInfo.TimeStamp, - esVersion: existingMem?.ESVersion ?? VersionInfo.UnknownVersion); + SetMember(members, newInfo, newInfo.Updated(newInfo.TimeStamp, + esVersion: existingMem?.ESVersion ?? VersionInfo.UnknownVersion)); } } } @@ -457,6 +459,22 @@ public static ClusterInfo MergeClusters(ClusterInfo myCluster, ClusterInfo other return new ClusterInfo(newMembers); } + private static bool IsSameMember(MemberInfo left, MemberInfo right) => + (left.InstanceId != Guid.Empty && right.InstanceId != Guid.Empty && left.InstanceId == right.InstanceId) || + left.Is(right.HttpEndPoint) || left.Is(right.ClusterEndPoint) || + right.Is(left.HttpEndPoint) || right.Is(left.ClusterEndPoint); + + private static void SetMember( + IDictionary members, + MemberInfo existing, + MemberInfo updated) + { + if (existing is not null) + members.Remove(existing.HttpEndPoint); + + members[updated.HttpEndPoint] = updated; + } + private static bool IsMoreUpToDate(MemberInfo member, MemberInfo existingMem) { if (member.EpochNumber != existingMem.EpochNumber) diff --git a/src/EventStore.Core/Services/Gossip/NodeGossipService.cs b/src/EventStore.Core/Services/Gossip/NodeGossipService.cs index cc25bfa63b..d917f18bbd 100644 --- a/src/EventStore.Core/Services/Gossip/NodeGossipService.cs +++ b/src/EventStore.Core/Services/Gossip/NodeGossipService.cs @@ -73,7 +73,7 @@ protected override MemberInfo GetInitialMe() lastEpoch == null ? Guid.Empty : lastEpoch.EpochId, _nodePriority, _memberInfo.IsReadOnlyReplica, _memberInfo.ESVersion, - _memberInfo.ReplicationEndPoint); + _memberInfo.ClusterEndPoint); } protected override MemberInfo GetUpdatedMe(MemberInfo me) diff --git a/src/EventStore.Core/Services/Replication/GrpcReplicaServiceSupervisor.cs b/src/EventStore.Core/Services/Replication/GrpcReplicaServiceSupervisor.cs index 52d2f29e3d..b4ca2b0724 100644 --- a/src/EventStore.Core/Services/Replication/GrpcReplicaServiceSupervisor.cs +++ b/src/EventStore.Core/Services/Replication/GrpcReplicaServiceSupervisor.cs @@ -277,7 +277,7 @@ private async ValueTask ReplaceActiveAsync( { active.Service = _factory.Create( new FencedPublisher(this, active), - new GrpcReplicaConnectionEndpoints(leader.HttpEndPoint, _advertisedReplicaEndPoint)); + new GrpcReplicaConnectionEndpoints(leader.ClusterEndPoint, _advertisedReplicaEndPoint)); SetActive(active); var task = active.Service.Start(); _trackTask(task); @@ -296,7 +296,8 @@ private async ValueTask ReplaceActiveAsync( { ClearActive(active); await StopAsync(active.Service); - Log.Warning(exception, "Failed to start replication stream to [{leaderEndPoint}].", leader.HttpEndPoint); + Log.Warning(exception, "Failed to start replication stream to [{leaderEndPoint}].", + leader.ClusterEndPoint); _publisher.Publish(new ReplicationMessage.LeaderConnectionFailed( leaderConnectionCorrelationId, leader)); } diff --git a/src/EventStore.Core/Services/Replication/ReplicationGrpcClient.cs b/src/EventStore.Core/Services/Replication/ReplicationGrpcClient.cs index 4c8f0ea9e9..463f63a119 100644 --- a/src/EventStore.Core/Services/Replication/ReplicationGrpcClient.cs +++ b/src/EventStore.Core/Services/Replication/ReplicationGrpcClient.cs @@ -1,6 +1,7 @@ using System; using System.Collections.Generic; using System.Net; +using System.Net.Http; using System.Runtime.CompilerServices; using System.Threading; using System.Threading.Tasks; @@ -32,24 +33,37 @@ public interface IReplicationGrpcCall : IDisposable public sealed class ReplicationGrpcClientFactory : IReplicationGrpcClientFactory { + private static readonly TimeSpan DefaultKeepAlivePingDelay = TimeSpan.FromMilliseconds(700); + private static readonly TimeSpan DefaultKeepAlivePingTimeout = TimeSpan.FromMilliseconds(700); + private static readonly TimeSpan MinimumKeepAlivePingValue = TimeSpan.FromSeconds(1); private readonly string _uriScheme; private readonly INodeHttpClientFactory _nodeHttpClientFactory; + private readonly TimeSpan _keepAlivePingDelay; + private readonly TimeSpan _keepAlivePingTimeout; public ReplicationGrpcClientFactory( string uriScheme, - INodeHttpClientFactory nodeHttpClientFactory) + INodeHttpClientFactory nodeHttpClientFactory, + TimeSpan? keepAlivePingDelay = null, + TimeSpan? keepAlivePingTimeout = null) { Ensure.NotNullOrEmpty(uriScheme, nameof(uriScheme)); Ensure.NotNull(nodeHttpClientFactory, nameof(nodeHttpClientFactory)); _uriScheme = uriScheme; _nodeHttpClientFactory = nodeHttpClientFactory; + _keepAlivePingDelay = NormalizeKeepAliveValue(keepAlivePingDelay ?? DefaultKeepAlivePingDelay); + _keepAlivePingTimeout = NormalizeKeepAliveValue(keepAlivePingTimeout ?? DefaultKeepAlivePingTimeout); } + private static TimeSpan NormalizeKeepAliveValue(TimeSpan value) => + value < MinimumKeepAlivePingValue ? MinimumKeepAlivePingValue : value; + public IReplicationGrpcClient Create(EndPoint leaderEndPoint) { Ensure.NotNull(leaderEndPoint, nameof(leaderEndPoint)); - return new ReplicationGrpcClient(_uriScheme, leaderEndPoint, _nodeHttpClientFactory); + return new ReplicationGrpcClient(_uriScheme, leaderEndPoint, _nodeHttpClientFactory, + _keepAlivePingDelay, _keepAlivePingTimeout); } } @@ -60,9 +74,18 @@ internal sealed class ReplicationGrpcClient : IReplicationGrpcClient public ReplicationGrpcClient( string uriScheme, EndPoint leaderEndPoint, - INodeHttpClientFactory nodeHttpClientFactory) + INodeHttpClientFactory nodeHttpClientFactory, + TimeSpan keepAlivePingDelay, + TimeSpan keepAlivePingTimeout) { - var httpClient = nodeHttpClientFactory.CreateHttpClient(leaderEndPoint.GetOtherNames()); + var httpClient = nodeHttpClientFactory.CreateHttpClient( + leaderEndPoint.GetOtherNames(), + handler => + { + handler.KeepAlivePingDelay = keepAlivePingDelay; + handler.KeepAlivePingTimeout = keepAlivePingTimeout; + handler.KeepAlivePingPolicy = HttpKeepAlivePingPolicy.Always; + }); httpClient.Timeout = Timeout.InfiniteTimeSpan; httpClient.DefaultRequestVersion = new Version(2, 0); diff --git a/src/EventStore.Core/Services/RequestForwarding/GrpcRequestForwardingSupervisor.cs b/src/EventStore.Core/Services/RequestForwarding/GrpcRequestForwardingSupervisor.cs index 0d91272c93..c8bb515477 100644 --- a/src/EventStore.Core/Services/RequestForwarding/GrpcRequestForwardingSupervisor.cs +++ b/src/EventStore.Core/Services/RequestForwarding/GrpcRequestForwardingSupervisor.cs @@ -214,11 +214,11 @@ private void ConnectToLeader(MemberInfo leader) IGrpcRequestForwardingService service = null; try { - var active = new ActiveStream(leader.InstanceId, leader.HttpEndPoint, connectionGeneration); + var active = new ActiveStream(leader.InstanceId, leader.ClusterEndPoint, connectionGeneration); service = _factory.Create( message => TryPublishIfActive(active, message), _publisher.Publish, - leader.HttpEndPoint, + leader.ClusterEndPoint, new ForwardingSessionGeneration(connectionGeneration)); active.Service = service; _active = active; @@ -237,7 +237,7 @@ private void ConnectToLeader(MemberInfo leader) service?.Stop(); Log.Warning(exception, "Failed to start request forwarding stream to [{leaderEndPoint}].", - leader.HttpEndPoint); + leader.ClusterEndPoint); if (connectionGeneration == _connectionGeneration) { ScheduleReconnect(leader.InstanceId, connectionGeneration); @@ -280,7 +280,7 @@ private bool HasHealthyStreamTo(MemberInfo leader) => _active is not null && !_active.Service.Task.IsCompleted && _active.LeaderId == leader.InstanceId && - Equals(_active.LeaderEndPoint, leader.HttpEndPoint); + Equals(_active.LeaderEndPoint, leader.ClusterEndPoint); private void PublishIfActive(ActiveStream active, Message message) { diff --git a/src/EventStore.Projections.Core.Tests/ClientAPI/Cluster/specification_with_standard_projections_runnning.cs b/src/EventStore.Projections.Core.Tests/ClientAPI/Cluster/specification_with_standard_projections_runnning.cs index 7e379c45be..44c8a71412 100644 --- a/src/EventStore.Projections.Core.Tests/ClientAPI/Cluster/specification_with_standard_projections_runnning.cs +++ b/src/EventStore.Projections.Core.Tests/ClientAPI/Cluster/specification_with_standard_projections_runnning.cs @@ -44,19 +44,19 @@ public abstract class specification_with_standard_projections_runnning Ports => _ports; @@ -79,7 +79,7 @@ public override async Task TestFixtureSetUp() { var gossipSeeds = _nodeEndpoints .Where((_, otherIndex) => otherIndex != index) - .Select(x => (EndPoint)x.HttpEndPoint) + .Select(x => (EndPoint)x.ClusterEndPoint) .ToArray(); _nodes[index] = CreateNode(index, _nodeEndpoints[index], gossipSeeds); } @@ -130,7 +130,7 @@ private MiniClusterNode CreateNode(int index, Endpoints e return new MiniClusterNode( PathName, index, - endpoints.InternalTcp, + endpoints.ClusterEndPoint, endpoints.ExternalTcp, endpoints.HttpEndPoint, subsystems: [_projections[index]],