Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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
Expand Down
2 changes: 1 addition & 1 deletion proto.lock
Original file line number Diff line number Diff line change
Expand Up @@ -8124,4 +8124,4 @@
}
}
]
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand All @@ -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);

Expand Down
Original file line number Diff line number Diff line change
@@ -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<EndpointBinding> _bindings;
private readonly IReadOnlyList<GrpcEndpointRoute> _routes;
private readonly EndpointRole _defaultRouteRole;
private readonly EndpointRole _nonIpEndpointRole;

public EndpointPolicy(
IReadOnlyList<EndpointBinding> bindings,
IReadOnlyList<GrpcEndpointRoute> 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);
}
}
51 changes: 45 additions & 6 deletions src/EventStore.ClusterNode/Program.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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 =>
{
Expand All @@ -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)
{
Expand Down Expand Up @@ -325,6 +349,16 @@ async Task Run(ClusterVNodeHostedService hostedService, ManualResetEventSlim sig
builder.Services.AddSingleton<IHostedService>(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<UiCredentialsMiddleware>();
hostedService.Node.Startup.Configure(app);
if (oauthEnabled)
Expand Down Expand Up @@ -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);
Expand Down
54 changes: 37 additions & 17 deletions src/EventStore.Core.Tests/Cluster/MemberInfoTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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];
Expand All @@ -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,
Expand All @@ -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,
Expand All @@ -121,28 +141,28 @@ 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)));
grpcCluster.Members[0].ReplicationEndPoint = null;

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) =>
Expand Down
Loading
Loading