Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
42 commits
Select commit Hold shift + click to select a range
37b7b1a
feat(cluster): preserve independent replication identity
yordis Sep 21, 2026
e322bb8
fix(cluster): preserve replication identity for monitoring
yordis Sep 21, 2026
9a1a9fc
fix(cluster): advertise the replication identity
yordis Sep 21, 2026
28f5203
feat(cluster): isolate internal replication traffic
yordis Sep 21, 2026
09439b0
fix(monitoring): match isolated replication members
yordis Sep 21, 2026
bcefb75
chore(transport): retire the internal TCP runtime
yordis Sep 21, 2026
2e91033
fix(tests): prove read-only delete rejection
yordis Sep 21, 2026
2e33e61
fix(tests): follow the gRPC replication endpoint
yordis Sep 21, 2026
43d6ce1
fix(tests): follow the gRPC membership matcher
yordis Sep 21, 2026
d79424b
fix(protocol): keep retired API out of the lockfile
yordis Sep 21, 2026
f14644b
feat(monitoring): preserve connection visibility over gRPC
yordis Sep 21, 2026
1a63389
fix(monitoring): preserve queue visibility
yordis Sep 21, 2026
834235c
fix(monitoring): report replication status independently
yordis Sep 21, 2026
c3dc51f
chore(transport): retire the legacy TCP package
yordis Sep 21, 2026
1c35ce7
fix(cluster): avoid drift from gRPC service contracts
yordis Sep 22, 2026
5b387fb
fix(tests): preserve read-only replica write coverage
yordis Sep 22, 2026
a80b6af
chore(transport): align with merged cluster endpoint contract
yordis Sep 22, 2026
fc98ac5
fix(grpc): prevent orphaned subscriptions on disconnect
yordis Sep 23, 2026
1e0b002
fix(grpc): preserve forwarded authentication failures
yordis Sep 23, 2026
d13926c
fix(monitoring): preserve connection visibility across endpoint changes
yordis Sep 23, 2026
fa68f7d
chore(transport): keep package retirement behind gRPC parity
yordis Sep 23, 2026
7430d58
fix(tests): keep forwarded authentication coverage in CI
yordis Sep 23, 2026
bfe81b0
fix(monitoring): retain verified authentication coverage
yordis Sep 23, 2026
403dc0a
chore(transport): retain verified authentication coverage
yordis Sep 23, 2026
bcb371f
chore(monitoring): guard against stale connection visibility
yordis Sep 23, 2026
3a340dd
chore(transport): keep connection visibility regression in the retire…
yordis Sep 23, 2026
0405597
chore(replication): protect secure cluster identity
yordis Sep 23, 2026
5c1281e
fix(monitoring): preserve connection visibility during queue failures
yordis Sep 23, 2026
ec97ff0
fix(monitoring): distinguish unavailable from independent connection …
yordis Sep 23, 2026
c08032d
chore(transport): preserve independent connection visibility
yordis Sep 23, 2026
8eadee8
fix(monitoring): restrict operational diagnostics to authorized callers
yordis Sep 23, 2026
8101a97
fix(container): avoid retired certificate dependency
yordis Sep 23, 2026
895ebaa
chore(transport): keep diagnostics protected through retirement
yordis Sep 23, 2026
df93d11
fix(monitoring): preserve independent replication diagnostics on queu…
yordis Sep 23, 2026
4975499
chore(transport): preserve diagnostics through retirement
yordis Sep 23, 2026
c3b1ff7
fix(monitoring): protect replication diagnostics with their own permi…
yordis Sep 23, 2026
6946ea6
chore(transport): preserve replication access boundaries
yordis Sep 23, 2026
1692c4f
chore(monitoring): align diagnostics with merged transport parity
yordis Sep 24, 2026
1de7dc3
chore(transport): keep package retirement aligned with merged parity
yordis Sep 24, 2026
7133fb0
fix(monitoring): preserve replication failure visibility
yordis Sep 24, 2026
028b784
chore(transport): keep retirement aligned with diagnostics
yordis Sep 24, 2026
5b77efa
chore(transport): align retirement with merged monitoring
yordis Sep 24, 2026
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
1 change: 0 additions & 1 deletion Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,6 @@ FROM mcr.microsoft.com/dotnet/sdk:10.0-${CONTAINER_RUNTIME} AS test
WORKDIR /build
COPY --from=build ./build/published-tests ./published-tests
COPY --from=build ./build/ci ./ci
COPY --from=build ./build/src/EventStore.Core.Tests/Services/Transport/Tcp/test_certificates/ca/ca.crt /usr/local/share/ca-certificates/ca_eventstore_test.crt
COPY ./scripts/test.sh /build/test.sh
RUN mkdir ./test-results
RUN chmod +x /build/test.sh
Expand Down
2 changes: 0 additions & 2 deletions scripts/test.sh
Original file line number Diff line number Diff line change
Expand Up @@ -3,8 +3,6 @@
set -o pipefail
set -o xtrace

update-ca-certificates

core_grpc_security_projects=(
EventStore.Core.Tests
)
Expand Down
1 change: 0 additions & 1 deletion src/Directory.Packages.props
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,6 @@
<PackageVersion Include="DotNext.IO" Version="5.26.3" />
<PackageVersion Include="DotNext.Threading" Version="5.26.3" />
<PackageVersion Include="DotNext.Unsafe" Version="5.26.2" />
<PackageVersion Include="EventStore.Client" Version="21.2.0" />
<PackageVersion Include="EventStore.Client.Grpc.Streams" Version="23.3.9" />
<PackageVersion Include="EventStore.Client.Grpc.Operations" Version="23.3.9" />
<!--Version 8 and beyond are free for open-source projects and non-commercial use, but commercial use requires a paid license-->
Expand Down
3 changes: 0 additions & 3 deletions src/EventStore.Core.Tests/EventStore.Core.Tests.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -27,11 +27,8 @@
<ProjectReference Include="..\EventStore.Core\EventStore.Core.csproj" />
<ProjectReference Include="..\EventStore.ClusterNode\EventStore.ClusterNode.csproj" PrivateAssets="all" />
<ProjectReference Include="..\EventStore.PluginHosting\EventStore.PluginHosting.csproj" />
<ProjectReference Include="..\EventStore.Transport.Tcp\EventStore.Transport.Tcp.csproj" />
Comment thread
cursor[bot] marked this conversation as resolved.
</ItemGroup>
<ItemGroup>
<EmbeddedResource Include="Services\Transport\Tcp\test_certificates\**\*.crt" />
<EmbeddedResource Include="Services\Transport\Tcp\test_certificates\**\*.key" />
<EmbeddedResource Remove="FakePlugin\**" />
</ItemGroup>
<ItemGroup>
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,144 @@
using System;
using System.Collections.Concurrent;
using System.Linq;
using System.Net;
using System.Net.Http;
using System.Security.Cryptography.X509Certificates;
using System.Threading.Tasks;
using EventStore.ClusterNode;
using EventStore.Common.Utils;
using EventStore.Core.Authorization;
using EventStore.Core.Bus;
using EventStore.Core.Messages;
using EventStore.Core.Messaging;
using EventStore.Core.Services.Replication;
using EventStore.Core.Services.Transport.Grpc.Replication;
using EventStore.Core.Services.Transport.Http.NodeHttpClientFactory;
using EventStore.Core.Tests.Helpers;
using Grpc.Core;
using Grpc.Net.Client;
using Microsoft.AspNetCore.Builder;
using Microsoft.AspNetCore.Hosting;
using Microsoft.AspNetCore.Hosting.Server;
using Microsoft.AspNetCore.Hosting.Server.Features;
using Microsoft.AspNetCore.Server.Kestrel.Core;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using NUnit.Framework;
using Proto = EventStore.Replication;

namespace EventStore.Core.Tests.Services.Transport.Grpc.Replication;

[TestFixture]
public class ReplicationMutualTlsTests
{
[Test]
public async Task trusted_client_certificate_reaches_replication_with_certificate_identity()
{
using var root = TestCertificates.GetRootCertificate();
using var serverCertificate = TestCertificates.GetServerCertificate();
using var clientCertificate = TestCertificates.GetOtherServerCertificate();
var publisher = new CapturingPublisher();
using var host = StartServer(serverCertificate, new X509Certificate2Collection(root), publisher);
using var channel = CreateChannel(host, clientCertificate, new X509Certificate2Collection(root));
var client = new Proto.Replication.ReplicationClient(channel);
using var call = client.Replicate(deadline: DateTime.UtcNow.AddSeconds(10));

await call.RequestStream.WriteAsync(SubscribeFrame());
await call.RequestStream.CompleteAsync();
Assert.That(await call.ResponseStream.MoveNext(), Is.False);
Assert.That(publisher.Messages.OfType<ReplicationMessage.ReplicaSubscriptionRequest>().Count(), Is.EqualTo(1));
Assert.That(publisher.Messages.OfType<ReplicationMessage.ReplicaSubscriptionRequest>().Single()
.Session.Identity.TransportIdentityKind,
Is.EqualTo(ReplicationTransportIdentityKind.ClientCertificateSha256));
}

[Test]
public async Task secure_replication_without_client_certificate_is_rejected_by_application_identity()
{
using var root = TestCertificates.GetRootCertificate();
using var serverCertificate = TestCertificates.GetServerCertificate();
var publisher = new CapturingPublisher();
using var host = StartServer(serverCertificate, new X509Certificate2Collection(root), publisher);
using var channel = CreateChannel(host, null, new X509Certificate2Collection(root));
var client = new Proto.Replication.ReplicationClient(channel);
using var call = client.Replicate(deadline: DateTime.UtcNow.AddSeconds(10));

await call.RequestStream.WriteAsync(SubscribeFrame());
await call.RequestStream.CompleteAsync();
var exception = Assert.ThrowsAsync<RpcException>(async () => await call.ResponseStream.MoveNext());
Assert.That(exception!.StatusCode, Is.EqualTo(StatusCode.Unauthenticated));
Assert.That(publisher.Messages, Is.Empty);
}

private static IHost StartServer(
X509Certificate2 serverCertificate,
X509Certificate2Collection trustedRoots,
IPublisher publisher)
{
var host = new HostBuilder()
.ConfigureWebHost(webHost => webHost
.UseKestrel(server => server.Listen(IPAddress.Loopback, 0, options =>
{
options.Protocols = HttpProtocols.Http2;
options.UseHttps(Program.CreateServerOptionsSelectionCallback(
() => serverCertificate,
() => null,
(certificate, chain, errors) => ClusterVNode<string>.ValidateClientCertificate(
certificate, chain, errors, () => null, () => trustedRoots)), null);
}))
.ConfigureServices(services =>
{
services.AddGrpc();
services.AddSingleton(new ReplicationService(
publisher, new PassthroughAuthorizationProvider()));
})
.Configure(app =>
{
app.UseRouting();
app.UseEndpoints(endpoints => endpoints.MapGrpcService<ReplicationService>());
}))
.Build();
host.Start();
return host;
}

private static GrpcChannel CreateChannel(
IHost host,
X509Certificate2 clientCertificate,
X509Certificate2Collection trustedRoots)
{
var factory = new NodeHttpClientFactory(
Uri.UriSchemeHttps,
(certificate, chain, errors, names) => ClusterVNode<string>.ValidateServerCertificate(
certificate, chain, errors, () => null, () => trustedRoots, names),
() => clientCertificate);
var httpClient = factory.CreateHttpClient(["localhost"]);
var address = host.Services.GetRequiredService<IServer>()
.Features.Get<IServerAddressesFeature>()!.Addresses.Single();
return GrpcChannel.ForAddress(address, new GrpcChannelOptions
{
HttpClient = httpClient,
DisposeHttpClient = true
});
}

private static Proto.ReplicaFrame SubscribeFrame() => ReplicationGrpcCodec.ToGrpc(
new ReplicationMessage.SubscribeReplica(
ReplicationSubscriptionVersions.V_CURRENT,
0,
Guid.NewGuid(),
[],
new DnsEndPoint("replica.internal", 1112),
Guid.NewGuid(),
Guid.NewGuid(),
true,
Guid.NewGuid()));

private sealed class CapturingPublisher : IPublisher
{
public ConcurrentQueue<Message> Messages { get; } = new();

public void Publish(Message message) => Messages.Enqueue(message);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,168 @@
using System;
using System.Linq;
using System.Net;
using System.Net.Http;
using System.Security.Cryptography.X509Certificates;
using System.Threading.Tasks;
using EventStore.ClusterNode;
using EventStore.Common.Utils;
using EventStore.Core.Services.Transport.Http.NodeHttpClientFactory;
using EventStore.Core.Tests.Helpers;
using Microsoft.AspNetCore.Builder;
using Microsoft.AspNetCore.Hosting;
using Microsoft.AspNetCore.Hosting.Server;
using Microsoft.AspNetCore.Hosting.Server.Features;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using NUnit.Framework;

namespace EventStore.Core.Tests.Services.Transport.Http;

[TestFixture]
public class ssl_connections_mutual_auth
{
[TestCase(true, true, true, true, true)]
[TestCase(true, false, true, true, false)]
[TestCase(false, true, true, true, false)]
[TestCase(false, false, true, true, false)]
[TestCase(true, true, true, false, true)]
[TestCase(true, false, true, false, true)]
[TestCase(false, true, true, false, false)]
[TestCase(false, false, true, false, false)]
[TestCase(true, true, false, true, true)]
[TestCase(true, false, false, true, false)]
[TestCase(false, true, false, true, true)]
[TestCase(false, false, false, true, false)]
[TestCase(true, true, false, false, true)]
[TestCase(true, false, false, false, true)]
[TestCase(false, true, false, false, true)]
[TestCase(false, false, false, false, true)]
public async Task connection_outcome_follows_server_and_client_certificate_policy(
bool useTrustedServerCertificate,
bool useTrustedClientCertificate,
bool validateServerCertificate,
bool validateClientCertificate,
bool shouldConnectSuccessfully)
{
using var rootCertificate = TestCertificates.GetRootCertificate();
using var serverCertificate = useTrustedServerCertificate
? TestCertificates.GetServerCertificate()
: TestCertificates.GetUntrustedCertificate();
using var clientCertificate = useTrustedClientCertificate
? TestCertificates.GetOtherServerCertificate()
: TestCertificates.GetUntrustedCertificate();
var trustedRoots = new X509Certificate2Collection(rootCertificate);
CertificateDelegates.ClientCertificateValidator clientValidator = validateClientCertificate
? (certificate, chain, errors) => ClusterVNode<string>.ValidateClientCertificate(
certificate,
chain,
errors,
() => null,
() => trustedRoots)
: (_, _, _) => (true, null);

using var host = StartServer(serverCertificate, clientValidator);
var connected = await TryConnect(
host,
clientCertificate,
validateServerCertificate,
trustedRoots);

Assert.That(connected, Is.EqualTo(shouldConnectSuccessfully));
}

[TestCase(true)]
[TestCase(false)]
public async Task client_certificate_is_optional_at_the_transport_boundary(bool validateServerCertificate)
{
using var rootCertificate = TestCertificates.GetRootCertificate();
using var serverCertificate = TestCertificates.GetServerCertificate();
var trustedRoots = new X509Certificate2Collection(rootCertificate);
using var host = StartServer(
serverCertificate,
(_, _, _) => throw new AssertionException("Missing client certificates bypass node validation."));

var connected = await TryConnect(
host,
clientCertificate: null,
validateServerCertificate,
trustedRoots);

Assert.That(connected, Is.True);
}

[Test]
public async Task server_certificate_is_required()
{
using var rootCertificate = TestCertificates.GetRootCertificate();
var trustedRoots = new X509Certificate2Collection(rootCertificate);
using var host = StartServer(
serverCertificate: null,
(_, _, _) => (true, null));

var connected = await TryConnect(
host,
clientCertificate: null,
validateServerCertificate: false,
trustedRoots);

Assert.That(connected, Is.False);
}

private static IHost StartServer(
X509Certificate2 serverCertificate,
CertificateDelegates.ClientCertificateValidator clientCertificateValidator)
{
var host = new HostBuilder()
.ConfigureWebHost(webHost => webHost
.UseKestrel(server => server.Listen(IPAddress.Loopback, 0, listenOptions =>
listenOptions.UseHttps(Program.CreateServerOptionsSelectionCallback(
() => serverCertificate,
() => null,
clientCertificateValidator), null)))
.Configure(app => app.Run(context => context.Response.CompleteAsync())))
.Build();
host.Start();
return host;
}

private static async Task<bool> TryConnect(
IHost host,
X509Certificate2 clientCertificate,
bool validateServerCertificate,
X509Certificate2Collection trustedRoots)
{
CertificateDelegates.ServerCertificateValidator serverValidator = validateServerCertificate
? (certificate, chain, errors, otherNames) => ClusterVNode<string>.ValidateServerCertificate(
certificate,
chain,
errors,
() => null,
() => trustedRoots,
otherNames)
: (_, _, _, _) => (true, null);
var clientFactory = new NodeHttpClientFactory(
Uri.UriSchemeHttps,
serverValidator,
() => clientCertificate);
using var client = clientFactory.CreateHttpClient(["localhost"]);
client.Timeout = TimeSpan.FromSeconds(5);
var address = host.Services.GetRequiredService<IServer>()
.Features.Get<IServerAddressesFeature>()!.Addresses.Single();
using var request = new HttpRequestMessage(HttpMethod.Get, address)
{
Version = HttpVersion.Version20,
VersionPolicy = HttpVersionPolicy.RequestVersionExact,
};

try
{
using var response = await client.SendAsync(request);
return response.IsSuccessStatusCode;
}
catch (HttpRequestException)
{
return false;
}
}
}
Loading
Loading