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
5 changes: 5 additions & 0 deletions proto.lock
Original file line number Diff line number Diff line change
Expand Up @@ -1988,6 +1988,11 @@
"id": 21,
"name": "es_version",
"type": "string"
},
{
"id": 22,
"name": "replication_end_point",
"type": "EndPoint"
}
]
}
Expand Down
105 changes: 105 additions & 0 deletions src/EventStore.Core.Tests/Cluster/MemberInfoTests.cs
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
using System;
using System.Net;
using System.Reflection;
using EventStore.Core.Data;
using NUnit.Framework;

Expand All @@ -8,6 +9,13 @@ namespace EventStore.Core.Tests.Cluster;
[TestFixture]
public class MemberInfoTests
{
private static readonly DnsEndPoint InternalTcp = new("internal", 1112);
private static readonly DnsEndPoint InternalSecureTcp = new("internal-secure", 2112);
private static readonly DnsEndPoint ExternalTcp = new("external", 1113);
private static readonly DnsEndPoint ExternalSecureTcp = new("external-secure", 2113);
private static readonly DnsEndPoint Http = new("http", 2113);
private static readonly DnsEndPoint Replication = new("replication", 3113);

[Test]
public void member_with_dns_endpoint_should_equal()
{
Expand Down Expand Up @@ -50,4 +58,101 @@ public void member_with_ip_endpoint_should_equal()
Assert.True(memberWithDnsEndPoint.Is(ipEndPoint));
Assert.True(memberWithDnsEndPoint.Is(dnsEndPoint));
}

[Test]
public void grpc_round_trip_preserves_tcp_and_replication_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);

var result = FromGrpcClusterInfo(ToGrpcClusterInfo(
new EventStore.Core.Cluster.ClusterInfo(member))).Members[0];

Assert.That(result.InternalTcpEndPoint, Is.EqualTo(InternalTcp));
Assert.That(result.InternalSecureTcpEndPoint, Is.Null);
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));
}

[Test]
public void explicit_replication_endpoint_is_recognized_without_replacing_tcp_endpoints()
{
var member = CreateMember(Replication);
var vnode = new VNodeInfo(Guid.NewGuid(), 0,
new IPEndPoint(IPAddress.Loopback, 1112), null,
new IPEndPoint(IPAddress.Loopback, 1113), null,
Http, false, Replication);
var advertise = new GossipAdvertiseInfo(
InternalTcp, InternalSecureTcp, ExternalTcp, ExternalSecureTcp, Http,
null, null, 0, null, 0, 0, Replication);

Assert.That(member.Is(Replication), Is.True);
Assert.That(member.InternalTcpEndPoint, Is.EqualTo(InternalTcp));
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));
}

[Test]
public void client_member_preserves_the_replication_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));
}

[Test]
public void missing_replication_endpoint_falls_back_to_http_endpoint()
{
var member = CreateMember();
var vnode = new VNodeInfo(Guid.NewGuid(), 0,
new IPEndPoint(IPAddress.Loopback, 1112), null,
new IPEndPoint(IPAddress.Loopback, 1113), null,
Http, false);
var advertise = new GossipAdvertiseInfo(
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));
}

[Test]
public void grpc_member_without_replication_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));
}

private static EventStore.Core.Cluster.MemberInfo CreateMember(DnsEndPoint replicationEndPoint = 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);

private static EventStore.Cluster.ClusterInfo ToGrpcClusterInfo(
EventStore.Core.Cluster.ClusterInfo clusterInfo) =>
(EventStore.Cluster.ClusterInfo)typeof(EventStore.Core.Cluster.ClusterInfo)
.GetMethod("ToGrpcClusterInfo", BindingFlags.NonPublic | BindingFlags.Static)!
.Invoke(null, [clusterInfo])!;

private static EventStore.Core.Cluster.ClusterInfo FromGrpcClusterInfo(
EventStore.Cluster.ClusterInfo clusterInfo) =>
(EventStore.Core.Cluster.ClusterInfo)typeof(EventStore.Core.Cluster.ClusterInfo)
.GetMethod("FromGrpcClusterInfo", BindingFlags.NonPublic | BindingFlags.Static)!
.Invoke(null, [clusterInfo, null])!;
}
Original file line number Diff line number Diff line change
Expand Up @@ -122,7 +122,8 @@ protected static MemberInfo MemberInfoForVNode(VNodeInfo nodeInfo, DateTime utcN
return MemberInfo.ForVNode(nodeInfo.InstanceId, utcNow, nodeState, isAlive,
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);
0, writerCheckpoint ?? 0, 0, -1, epochNumber ?? -1, Guid.Empty, nodePriority ?? 0, false, esVersion,
nodeInfo.ReplicationEndPoint);
}

/// <summary>
Expand Down Expand Up @@ -245,6 +246,42 @@ public void should_start_gossiping_and_schedule_another_gossip()
}
}

public class when_got_gossip_seed_sources_with_distinct_replication_endpoint : NodeGossipServiceTestFixture
{
public when_got_gossip_seed_sources_with_distinct_replication_endpoint()
{
_currentNode = new VNodeInfo(
Guid.Parse("00000000-0000-0000-0000-000000000001"), 1,
new IPEndPoint(IPAddress.Loopback, 1111),
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, 1112));
}

protected override Message[] Given() => [new SystemMessage.SystemInit()];

protected override Message When() =>
new GossipMessage.GotGossipSeedSources([
_currentNode.HttpEndPoint,
_nodeTwo.HttpEndPoint,
_nodeThree.HttpEndPoint
]);

[Test]
public void should_preserve_the_replication_endpoint()
{
var gossip = (GossipMessage.SendGossip)_bus.Messages
.OfType<GrpcMessage.SendOverGrpc>()
.Single()
.Message;
var currentMember = gossip.ClusterInfo.Members.Single(x => x.InstanceId == _currentNode.InstanceId);

Assert.That(currentMember.ReplicationEndPoint, Is.EqualTo(_currentNode.ReplicationEndPoint));
}
}

public class when_gossip : NodeGossipServiceTestFixture
{
private int _gossipRound = GossipServiceBase.GossipRoundStartupThreshold + 1;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -151,7 +151,12 @@ options with
[Test]
public void should_advertise_the_configured_replication_port()
{
Assert.AreEqual(3112, _node.GossipAdvertiseInfo.InternalSecureTcp.Port);
Assert.Multiple(() =>
{
Assert.AreEqual(_options.Interface.ReplicationPort, _node.NodeInfo.ReplicationEndPoint.GetPort());
Assert.AreEqual(3112, _node.GossipAdvertiseInfo.ReplicationEndPoint.Port);
Assert.AreEqual(3112, _node.GossipAdvertiseInfo.InternalSecureTcp.Port);
});
}
}

Expand Down
6 changes: 5 additions & 1 deletion src/EventStore.Core/Cluster/ClientClusterInfo.cs
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,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 HttpEndPointIp { get; set; }
public int HttpEndPointPort { get; set; }
Expand Down Expand Up @@ -84,6 +86,8 @@ public ClientMemberInfo(MemberInfo member)

InternalHttpEndPointIp = member.HttpEndPoint.GetHost();
InternalHttpEndPointPort = member.HttpEndPoint.GetPort();
ReplicationEndPointIp = member.ReplicationEndPoint.GetHost();
ReplicationEndPointPort = member.ReplicationEndPoint.GetPort();

HttpEndPointIp = string.IsNullOrEmpty(member.AdvertiseHostToClientAs)
? member.HttpEndPoint.GetHost()
Expand Down Expand Up @@ -125,6 +129,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}, " +
$"HttpEndPointIp: {HttpEndPointIp}, HttpEndPointPort: {HttpEndPointPort}, " +
$"LastCommitPosition: {LastCommitPosition}, WriterCheckpoint: {WriterCheckpoint}, ChaserCheckpoint: {ChaserCheckpoint}, " +
$"EpochPosition: {EpochPosition}, EpochNumber: {EpochNumber}, EpochId: {EpochId:B}, NodePriority: {NodePriority}, " +
Expand All @@ -133,4 +138,3 @@ public override string ToString()
}
}
}

9 changes: 8 additions & 1 deletion src/EventStore.Core/Cluster/ClusterInfo.cs
Original file line number Diff line number Diff line change
Expand Up @@ -83,7 +83,11 @@ internal static ClusterInfo FromGrpcClusterInfo(EventStore.Cluster.ClusterInfo g
x.AdvertiseHostToClientAs, (int)x.AdvertiseHttpPortToClientAs, (int)x.AdvertiseTcpPortToClientAs,
x.LastCommitPosition, x.WriterCheckpoint, x.ChaserCheckpoint,
x.EpochPosition, x.EpochNumber, Uuid.FromDto(x.EpochId).ToGuid(), x.NodePriority,
x.IsReadOnlyReplica, x.EsVersion == String.Empty ? null : x.EsVersion
x.IsReadOnlyReplica, x.EsVersion == String.Empty ? null : x.EsVersion,
x.ReplicationEndPoint is null
? null
: new DnsEndPoint(x.ReplicationEndPoint.Address, (int)x.ReplicationEndPoint.Port)
.WithClusterDns(clusterDns)
)).ToArray();
return new ClusterInfo(receivedMembers);
}
Expand All @@ -99,6 +103,9 @@ internal static EventStore.Cluster.ClusterInfo ToGrpcClusterInfo(ClusterInfo clu
HttpEndPoint = new EventStore.Cluster.EndPoint(
x.HttpEndPoint.GetHost(),
(uint)x.HttpEndPoint.GetPort()),
ReplicationEndPoint = new EventStore.Cluster.EndPoint(
x.ReplicationEndPoint.GetHost(),
(uint)x.ReplicationEndPoint.GetPort()),
InternalTcp = x.InternalSecureTcpEndPoint != null ?
new EventStore.Cluster.EndPoint(
x.InternalSecureTcpEndPoint.GetHost(),
Expand Down
33 changes: 23 additions & 10 deletions src/EventStore.Core/Cluster/MemberInfo.cs
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ public class MemberInfo : IEquatable<MemberInfo>
public readonly EndPoint ExternalTcpEndPoint;
public readonly EndPoint ExternalSecureTcpEndPoint;
public readonly EndPoint HttpEndPoint;
public readonly EndPoint ReplicationEndPoint;
public readonly string AdvertiseHostToClientAs;
public readonly int AdvertiseHttpPortToClientAs;
public readonly int AdvertiseTcpPortToClientAs;
Expand All @@ -37,12 +38,13 @@ public class MemberInfo : IEquatable<MemberInfo>
public readonly string ESVersion;

public static MemberInfo ForManager(Guid instanceId, DateTime timeStamp, bool isAlive,
EndPoint httpEndPoint, string esVersion = VersionInfo.UnknownVersion)
EndPoint httpEndPoint, string esVersion = VersionInfo.UnknownVersion,
EndPoint replicationEndPoint = 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);
-1, -1, -1, -1, -1, Guid.Empty, 0, false, esVersion, replicationEndPoint);
}

public static MemberInfo ForVNode(Guid instanceId,
Expand All @@ -64,7 +66,8 @@ public static MemberInfo ForVNode(Guid instanceId,
int epochNumber,
Guid epochId,
int nodePriority,
bool isReadOnlyReplica, string esVersion = VersionInfo.UnknownVersion)
bool isReadOnlyReplica, string esVersion = VersionInfo.UnknownVersion,
EndPoint replicationEndPoint = null)
{
if (state == VNodeState.Manager)
{
Expand All @@ -76,7 +79,8 @@ public static MemberInfo ForVNode(Guid instanceId,
externalTcpEndPoint, externalSecureTcpEndPoint,
httpEndPoint, advertiseHostToClientAs, advertiseHttpPortToClientAs, advertiseTcpPortToClientAs,
lastCommitPosition, writerCheckpoint, chaserCheckpoint,
epochPosition, epochNumber, epochId, nodePriority, isReadOnlyReplica, esVersion);
epochPosition, epochNumber, epochId, nodePriority, isReadOnlyReplica, esVersion,
replicationEndPoint);
}

public static MemberInfo Initial(Guid instanceId,
Expand All @@ -92,7 +96,8 @@ public static MemberInfo Initial(Guid instanceId,
int advertiseHttpPortToClientAs,
int advertiseTcpPortToClientAs,
int nodePriority,
bool isReadOnlyReplica, string esVersion = VersionInfo.UnknownVersion)
bool isReadOnlyReplica, string esVersion = VersionInfo.UnknownVersion,
EndPoint replicationEndPoint = null)
{
if (state == VNodeState.Manager)
{
Expand All @@ -103,15 +108,17 @@ public static MemberInfo Initial(Guid instanceId,
internalTcpEndPoint, internalSecureTcpEndPoint,
externalTcpEndPoint, externalSecureTcpEndPoint,
httpEndPoint, advertiseHostToClientAs, advertiseHttpPortToClientAs, advertiseTcpPortToClientAs,
-1, -1, -1, -1, -1, Guid.Empty, nodePriority, isReadOnlyReplica, esVersion);
-1, -1, -1, -1, -1, Guid.Empty, nodePriority, isReadOnlyReplica, esVersion,
replicationEndPoint);
}

internal MemberInfo(Guid instanceId, DateTime timeStamp, VNodeState state, bool isAlive,
EndPoint internalTcpEndPoint, EndPoint internalSecureTcpEndPoint,
EndPoint externalTcpEndPoint, EndPoint externalSecureTcpEndPoint,
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)
long epochPosition, int epochNumber, Guid epochId, int nodePriority, bool isReadOnlyReplica,
string esVersion = null, EndPoint replicationEndPoint = null)
{
Ensure.Equal(false, internalTcpEndPoint == null && internalSecureTcpEndPoint == null, "Both internal TCP endpoints are null");
Ensure.NotNull(httpEndPoint, nameof(httpEndPoint));
Expand All @@ -127,6 +134,7 @@ internal MemberInfo(Guid instanceId, DateTime timeStamp, VNodeState state, bool
ExternalTcpEndPoint = externalTcpEndPoint;
ExternalSecureTcpEndPoint = externalSecureTcpEndPoint;
HttpEndPoint = httpEndPoint;
ReplicationEndPoint = replicationEndPoint ?? httpEndPoint;
AdvertiseHostToClientAs = advertiseHostToClientAs;
AdvertiseHttpPortToClientAs = advertiseHttpPortToClientAs;
AdvertiseTcpPortToClientAs = advertiseTcpPortToClientAs;
Expand Down Expand Up @@ -160,6 +168,7 @@ internal MemberInfo(MemberInfoDto dto)
? new DnsEndPoint(dto.ExternalTcpIp, dto.ExternalSecureTcpPort)
: null;
HttpEndPoint = new DnsEndPoint(dto.HttpEndPointIp, dto.HttpEndPointPort);
ReplicationEndPoint = HttpEndPoint;
AdvertiseHostToClientAs = dto.AdvertiseHostToClientAs;
AdvertiseHttpPortToClientAs = dto.AdvertiseHttpPortToClientAs;
AdvertiseTcpPortToClientAs = dto.AdvertiseTcpPortToClientAs;
Expand All @@ -176,11 +185,12 @@ internal MemberInfo(MemberInfoDto dto)
public bool Is(EndPoint endPoint)
{
return endPoint != null
&& HttpEndPoint.EndPointEquals(endPoint)
&& (HttpEndPoint.EndPointEquals(endPoint)
|| ReplicationEndPoint.EndPointEquals(endPoint)
|| (InternalTcpEndPoint != null && InternalTcpEndPoint.EndPointEquals(endPoint))
|| (InternalSecureTcpEndPoint != null && InternalSecureTcpEndPoint.EndPointEquals(endPoint))
|| (ExternalTcpEndPoint != null && ExternalTcpEndPoint.EndPointEquals(endPoint))
|| (ExternalSecureTcpEndPoint != null && ExternalSecureTcpEndPoint.EndPointEquals(endPoint));
|| (ExternalSecureTcpEndPoint != null && ExternalSecureTcpEndPoint.EndPointEquals(endPoint)));
}

public MemberInfo Updated(DateTime utcNow,
Expand Down Expand Up @@ -211,7 +221,7 @@ public MemberInfo Updated(DateTime utcNow,
epoch != null ? epoch.EpochNumber : EpochNumber,
epoch != null ? epoch.EpochId : EpochId,
nodePriority ?? NodePriority,
IsReadOnlyReplica, esVersion ?? ESVersion);
IsReadOnlyReplica, esVersion ?? ESVersion, ReplicationEndPoint);
}

public override string ToString()
Expand All @@ -228,6 +238,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}, " +
$"{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}";
Expand All @@ -254,6 +265,7 @@ public bool Equals(MemberInfo other)
&& Equals(other.ExternalTcpEndPoint, ExternalTcpEndPoint)
&& Equals(other.ExternalSecureTcpEndPoint, ExternalSecureTcpEndPoint)
&& Equals(other.HttpEndPoint, HttpEndPoint)
&& Equals(other.ReplicationEndPoint, ReplicationEndPoint)
&& other.AdvertiseHostToClientAs == AdvertiseHostToClientAs
&& other.AdvertiseHttpPortToClientAs == AdvertiseHttpPortToClientAs
&& other.AdvertiseTcpPortToClientAs == AdvertiseTcpPortToClientAs
Expand Down Expand Up @@ -299,6 +311,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) ^ (AdvertiseHostToClientAs != null ? AdvertiseHostToClientAs.GetHashCode() : 0);
result = (result * 397) ^ AdvertiseHttpPortToClientAs.GetHashCode();
result = (result * 397) ^ AdvertiseTcpPortToClientAs.GetHashCode();
Expand Down
Loading
Loading