diff --git a/src/EventStore.ClusterNode/AssemblyInfo.cs b/src/EventStore.ClusterNode/AssemblyInfo.cs new file mode 100644 index 0000000000..e240e9f65e --- /dev/null +++ b/src/EventStore.ClusterNode/AssemblyInfo.cs @@ -0,0 +1,3 @@ +using System.Runtime.CompilerServices; + +[assembly: InternalsVisibleTo("EventStore.Core.Tests")] diff --git a/src/EventStore.ClusterNode/Program.cs b/src/EventStore.ClusterNode/Program.cs index 270308e3c5..a82abd59cb 100644 --- a/src/EventStore.ClusterNode/Program.cs +++ b/src/EventStore.ClusterNode/Program.cs @@ -433,15 +433,24 @@ private static void TryListenOnUnixSocket(ClusterVNodeHostedService hostedServic } private static ServerOptionsSelectionCallback CreateServerOptionsSelectionCallback( - ClusterVNodeHostedService hostedService) + ClusterVNodeHostedService hostedService) => + CreateServerOptionsSelectionCallback( + hostedService.Node.CertificateSelector, + hostedService.Node.IntermediateCertificatesSelector, + hostedService.Node.InternalClientCertificateValidator); + + internal static ServerOptionsSelectionCallback CreateServerOptionsSelectionCallback( + Func certificateSelector, + Func intermediateCertificatesSelector, + CertificateDelegates.ClientCertificateValidator clientCertificateValidator) { return ((_, _, _, _) => { var serverOptions = new SslServerAuthenticationOptions { ServerCertificateContext = SslStreamCertificateContext.Create( - hostedService.Node.CertificateSelector(), - hostedService.Node.IntermediateCertificatesSelector(), + certificateSelector(), + intermediateCertificatesSelector(), offline: true), ClientCertificateRequired = true, // request a client certificate but it's not necessary for the client to supply one @@ -452,11 +461,10 @@ private static ServerOptionsSelectionCallback CreateServerOptionsSelectionCallba return true; } - var (isValid, error) = - hostedService.Node.InternalClientCertificateValidator( - certificate, - chain, - sslPolicyErrors); + var (isValid, error) = clientCertificateValidator( + certificate, + chain, + sslPolicyErrors); if (!isValid && error != null) { Log.Error("Client certificate validation error: {e}", error); diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/OperationsTests/AdminTests.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/OperationsTests/AdminTests.cs index 0c4cf4e782..b191ca6daf 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Grpc/OperationsTests/AdminTests.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/OperationsTests/AdminTests.cs @@ -268,6 +268,63 @@ public void returns_permission_denied() } } + [TestFixture(typeof(LogFormat.V2), typeof(string))] + public class when_starting_scavenge_as_admin : GrpcSpecification + { + private ScavengeResp _response; + + protected override Task Given() => Task.CompletedTask; + + protected override async Task When() + { + _response = await Channel.CreateCallInvoker().AsyncUnaryCall( + StartScavengeMethod, + null, + GetCallOptions(AdminCredentials), + new StartScavengeReq()); + } + + [Test] + public void returns_started_with_a_scavenge_id() + { + Assert.AreEqual(ScavengeResp.Types.ScavengeResult.Started, _response.ScavengeResult); + Assert.IsNotEmpty(_response.ScavengeId); + } + } + + [TestFixture(typeof(LogFormat.V2), typeof(string))] + public class when_starting_scavenge_while_one_is_running + : GrpcSpecification + { + private ScavengeResp _startedResponse; + private ScavengeResp _inProgressResponse; + + protected override Task Given() => Task.CompletedTask; + + protected override async Task When() + { + _startedResponse = await Channel.CreateCallInvoker().AsyncUnaryCall( + StartScavengeMethod, + null, + GetCallOptions(AdminCredentials), + new StartScavengeReq()); + + _inProgressResponse = await Channel.CreateCallInvoker().AsyncUnaryCall( + StartScavengeMethod, + null, + GetCallOptions(AdminCredentials), + new StartScavengeReq()); + } + + [Test] + public void returns_in_progress_with_the_running_scavenge_id() + { + Assert.AreEqual(ScavengeResp.Types.ScavengeResult.Started, _startedResponse.ScavengeResult); + Assert.AreEqual(ScavengeResp.Types.ScavengeResult.InProgress, _inProgressResponse.ScavengeResult); + Assert.AreEqual(_startedResponse.ScavengeId, _inProgressResponse.ScavengeId); + } + } + [TestFixture(typeof(LogFormat.V2), typeof(string))] public class when_starting_scavenge_without_permissions : GrpcSpecification { diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/read_all_events_forward_with_hard_deleted_stream_should.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/read_all_events_forward_with_hard_deleted_stream_should.cs index c476f5ab42..e61f0b8eb4 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/read_all_events_forward_with_hard_deleted_stream_should.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/read_all_events_forward_with_hard_deleted_stream_should.cs @@ -17,6 +17,7 @@ public class read_all_events_forward_with_hard_deleted_stream_should); private readonly List _allEvents = new(); + private readonly List _allEventsBackward = new(); private RpcException _streamReadException; protected override async Task Given() @@ -56,10 +57,17 @@ protected override async Task When() _streamReadException = ex; } - using var allCall = StreamsClient.Read(ReadAllRequest(), GetCallOptions(AdminCredentials)); + using var allCall = StreamsClient.Read(ReadAllRequest( + ReadReq.Types.Options.Types.ReadDirection.Forwards), GetCallOptions(AdminCredentials)); _allEvents.AddRange((await allCall.ResponseStream.ReadAllAsync().ToArrayAsync()) .Where(x => x.ContentCase == ReadResp.ContentOneofCase.Event) .Select(x => x.Event)); + + using var allBackwardCall = StreamsClient.Read(ReadAllRequest( + ReadReq.Types.Options.Types.ReadDirection.Backwards), GetCallOptions(AdminCredentials)); + _allEventsBackward.AddRange((await allBackwardCall.ResponseStream.ReadAllAsync().ToArrayAsync()) + .Where(x => x.ContentCase == ReadResp.ContentOneofCase.Event) + .Select(x => x.Event)); } [Test] @@ -82,6 +90,20 @@ public void returns_all_events_including_tombstone() Assert.That(streamEvents.Take(20).All(x => x.Event.Metadata[GrpcConstants.Metadata.Type] == "-"), Is.True); Assert.That(streamEvents[^1].Event.Metadata[GrpcConstants.Metadata.Type], Is.EqualTo(SystemEventTypes.StreamDeleted)); + Assert.That(streamEvents[^1].Event.StreamRevision, Is.EqualTo((ulong)long.MaxValue)); + } + + [Test] + public void returns_tombstone_from_backward_read_without_downgrading_its_revision() + { + var streamEvents = _allEventsBackward + .Where(x => x.Event.StreamIdentifier.StreamName.ToStringUtf8() == StreamName) + .ToArray(); + + Assert.That(streamEvents, Has.Length.EqualTo(21)); + Assert.That(streamEvents[0].Event.Metadata[GrpcConstants.Metadata.Type], + Is.EqualTo(SystemEventTypes.StreamDeleted)); + Assert.That(streamEvents[0].Event.StreamRevision, Is.EqualTo((ulong)long.MaxValue)); } private static ReadReq ReadStreamRequest() => new() @@ -100,15 +122,17 @@ public void returns_all_events_including_tombstone() } }; - private static ReadReq ReadAllRequest() => new() + private static ReadReq ReadAllRequest(ReadReq.Types.Options.Types.ReadDirection direction) => new() { Options = new() { UuidOption = new() { Structured = new() }, NoFilter = new(), - ReadDirection = ReadReq.Types.Options.Types.ReadDirection.Forwards, + ReadDirection = direction, Count = 100, - All = new() { Start = new() } + All = direction == ReadReq.Types.Options.Types.ReadDirection.Forwards + ? new() { Start = new() } + : new() { End = new() } } }; } diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/subscribe_should.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/subscribe_should.cs index 6c990134e0..594e916c14 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/subscribe_should.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/subscribe_should.cs @@ -62,5 +62,6 @@ await StreamsClient.TombstoneAsync(new() Assert.That(deleted, Is.Not.Null); Assert.That(deleted.Event.Event.StreamIdentifier.StreamName.ToStringUtf8(), Is.EqualTo(streamName)); Assert.That(deleted.Event.Event.Metadata[GrpcMetadata.Type], Is.EqualTo(SystemEventTypes.StreamDeleted)); + Assert.That(deleted.Event.Event.StreamRevision, Is.EqualTo((ulong)long.MaxValue)); } } diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/subscribe_to_all_should.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/subscribe_to_all_should.cs index 90d5b20462..994da05139 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/subscribe_to_all_should.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/subscribe_to_all_should.cs @@ -66,6 +66,7 @@ await StreamsClient.TombstoneAsync(new() var deleted = await ReadNextEventResponse(subscription); Assert.That(deleted.Event.StreamIdentifier.StreamName.ToStringUtf8(), Is.EqualTo(streamName)); Assert.That(deleted.Event.Metadata[GrpcMetadata.Type], Is.EqualTo(SystemEventTypes.StreamDeleted)); + Assert.That(deleted.Event.StreamRevision, Is.EqualTo((ulong)long.MaxValue)); } private static ReadReq SubscribeRequest() => new() diff --git a/src/EventStore.Core.Tests/Services/Transport/Http/with_intermediate_certificates.cs b/src/EventStore.Core.Tests/Services/Transport/Http/with_intermediate_certificates.cs new file mode 100644 index 0000000000..10a99612bb --- /dev/null +++ b/src/EventStore.Core.Tests/Services/Transport/Http/with_intermediate_certificates.cs @@ -0,0 +1,74 @@ +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.Tests.Certificates; +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 with_intermediate_certificates : with_certificate_chain_of_length_3 +{ + private IHost _host; + + [SetUp] + public void SetUp() + { + var certificate = X509CertificateLoader.LoadPkcs12(_leaf.ExportToPkcs12(), null); + _host = new HostBuilder() + .ConfigureWebHost(webHost => webHost + .UseKestrel(server => server.Listen(IPAddress.Loopback, 0, listenOptions => + listenOptions.UseHttps(Program.CreateServerOptionsSelectionCallback( + () => certificate, + () => new X509Certificate2Collection(_intermediate), + (_, _, _) => (true, null)), null))) + .Configure(app => app.Run(context => context.Response.CompleteAsync()))) + .Build(); + _host.Start(); + } + + [Test] + public async Task server_should_send_intermediate_certificate_during_handshake() + { + var handler = new SocketsHttpHandler(); + var gotLeaf = false; + var gotIntermediate = false; + handler.SslOptions.RemoteCertificateValidationCallback = (_, certificate, chain, _) => + { + gotLeaf = certificate is not null && certificate.GetCertHashString() == _leaf.GetCertHashString(); + gotIntermediate = chain is not null && chain.ChainElements.Cast() + .Any(element => element.Certificate.Thumbprint == _intermediate.Thumbprint); + return true; + }; + using var client = new HttpClient(handler); + var address = _host.Services.GetRequiredService() + .Features.Get()!.Addresses.Single(); + using var request = new HttpRequestMessage(HttpMethod.Get, address) + { + Version = HttpVersion.Version20, + VersionPolicy = HttpVersionPolicy.RequestVersionExact, + }; + + using var response = await client.SendAsync(request); + + Assert.That(gotLeaf, Is.True); + Assert.That(gotIntermediate, Is.True); + } + + [TearDown] + public void TearDown() + { + _host?.Dispose(); + } +} diff --git a/src/EventStore.Core.Tests/Services/Transport/Tcp/TcpClientDispatcherTests.cs b/src/EventStore.Core.Tests/Services/Transport/Tcp/TcpClientDispatcherTests.cs deleted file mode 100644 index b55b4d7d63..0000000000 --- a/src/EventStore.Core.Tests/Services/Transport/Tcp/TcpClientDispatcherTests.cs +++ /dev/null @@ -1,317 +0,0 @@ -using System; -using System.Collections.Generic; -using System.Linq; -using System.Text; -using EventStore.Client.Messages; -using EventStore.Core.Authentication.InternalAuthentication; -using EventStore.Core.Bus; -using EventStore.Core.Data; -using EventStore.Core.LogV2; -using EventStore.Core.Messages; -using EventStore.Core.Messaging; -using EventStore.Core.Services; -using EventStore.Core.Services.Transport.Tcp; -using EventStore.Core.Services.UserManagement; -using EventStore.Core.Tests.Authentication; -using EventStore.Core.Tests.Authorization; -using EventStore.Core.TransactionLog.LogRecords; -using EventStore.Core.Util; -using NUnit.Framework; -using EventRecord = EventStore.Core.Data.EventRecord; -using ResolvedEvent = EventStore.Core.Data.ResolvedEvent; - -namespace EventStore.Core.Tests.Services.Transport.Tcp; - -[TestFixture] -public class TcpClientDispatcherTests -{ - private readonly NoopEnvelope _envelope = new NoopEnvelope(); - - private ClientTcpDispatcher _dispatcher; - private TcpConnectionManager _connection; - - [OneTimeSetUp] - public void Setup() - { - _dispatcher = new ClientTcpDispatcher(2000); - - var dummyConnection = new DummyTcpConnection(); - _connection = new TcpConnectionManager( - Guid.NewGuid().ToString(), TcpServiceType.External, new ClientTcpDispatcher(2000), - new SynchronousScheduler(), dummyConnection, new SynchronousScheduler(), new InternalAuthenticationProvider( - InMemoryBus.CreateTest(), new Core.Helpers.IODispatcher(new SynchronousScheduler(), new NoopEnvelope()), - new StubPasswordHashAlgorithm(), 1, false, DefaultData.DefaultUserOptions), - new AuthorizationGateway(new TestAuthorizationProvider()), - TimeSpan.FromSeconds(10), TimeSpan.FromSeconds(10), (man, err) => { }, - Opts.ConnectionPendingSendBytesThresholdDefault, Opts.ConnectionQueueSizeThresholdDefault); - } - - [Test] - public void - when_wrapping_read_stream_events_forward_and_stream_was_deleted_should_not_downgrade_last_event_number_for_v2_clients() - { - var msg = new ClientMessage.ReadStreamEventsForwardCompleted(Guid.NewGuid(), "test-stream", 0, 100, - ReadStreamResult.StreamDeleted, new ResolvedEvent[0], new StreamMetadata(), - true, "", -1, long.MaxValue, true, 1000); - - var package = _dispatcher.WrapMessage(msg, (byte)ClientVersion.V2); - Assert.IsNotNull(package, "Package is null"); - Assert.AreEqual(TcpCommand.ReadStreamEventsForwardCompleted, package.Value.Command, "TcpCommand"); - - var dto = package.Value.Data.Deserialize(); - Assert.IsNotNull(dto, "DTO is null"); - - Assert.AreEqual(long.MaxValue, dto.LastEventNumber, "Last Event Number"); - } - - [Test] - public void - when_wrapping_read_stream_events_backward_and_stream_was_deleted_should_not_downgrade_last_event_number_for_v2_clients() - { - var msg = new ClientMessage.ReadStreamEventsBackwardCompleted(Guid.NewGuid(), "test-stream", 0, 100, - ReadStreamResult.StreamDeleted, new ResolvedEvent[0], new StreamMetadata(), - true, "", -1, long.MaxValue, true, 1000); - - var package = _dispatcher.WrapMessage(msg, (byte)ClientVersion.V2); - Assert.IsNotNull(package, "Package is null"); - Assert.AreEqual(TcpCommand.ReadStreamEventsBackwardCompleted, package.Value.Command, "TcpCommand"); - - var dto = package.Value.Data.Deserialize(); - Assert.IsNotNull(dto, "DTO is null"); - - Assert.AreEqual(long.MaxValue, dto.LastEventNumber, "Last Event Number"); - } - - [Test] - public void - when_wrapping_read_all_events_forward_completed_with_deleted_event_should_not_downgrade_last_event_number_for_v2_clients() - { - var events = new ResolvedEvent[] { - ResolvedEvent.ForUnresolvedEvent(CreateDeletedEventRecord(), 0), - }; - var msg = new ClientMessage.ReadAllEventsForwardCompleted(Guid.NewGuid(), ReadAllResult.Success, "", events, - new StreamMetadata(), true, 10, new TFPos(0, 0), - new TFPos(200, 200), new TFPos(0, 0), 100); - - var package = _dispatcher.WrapMessage(msg, (byte)ClientVersion.V2); - Assert.IsNotNull(package, "Package is null"); - Assert.AreEqual(TcpCommand.ReadAllEventsForwardCompleted, package.Value.Command, "TcpCommand"); - - var dto = package.Value.Data.Deserialize(); - Assert.IsNotNull(dto, "DTO is null"); - Assert.AreEqual(1, dto.Events.Count(), "Number of events"); - - Assert.AreEqual(long.MaxValue, dto.Events[0].Event.EventNumber, "Event Number"); - } - - [Test] - public void - when_wrapping_read_all_events_forward_completed_with_link_to_deleted_event_should_not_downgrade_version_for_v2_clients() - { - var events = new ResolvedEvent[] { - ResolvedEvent.ForResolvedLink(CreateLinkEventRecord(), CreateDeletedEventRecord(), 100) - }; - var msg = new ClientMessage.ReadAllEventsForwardCompleted(Guid.NewGuid(), ReadAllResult.Success, "", events, - new StreamMetadata(), true, 10, new TFPos(0, 0), - new TFPos(200, 200), new TFPos(0, 0), 100); - - var package = _dispatcher.WrapMessage(msg, (byte)ClientVersion.V2); - Assert.IsNotNull(package, "Package is null"); - Assert.AreEqual(TcpCommand.ReadAllEventsForwardCompleted, package.Value.Command, "TcpCommand"); - - var dto = package.Value.Data.Deserialize(); - Assert.IsNotNull(dto, "DTO is null"); - Assert.AreEqual(1, dto.Events.Count(), "Number of events"); - - Assert.AreEqual(0, dto.Events[0].Event.EventNumber, "Event Number"); - Assert.AreEqual(long.MaxValue, dto.Events[0].Link.EventNumber, "Link Event Number"); - } - - [Test] - public void - when_wrapping_read_all_events_backward_completed_with_deleted_event_should_not_downgrade_version_for_v2_clients() - { - var events = new ResolvedEvent[] { - ResolvedEvent.ForUnresolvedEvent(CreateDeletedEventRecord(), 0), - }; - var msg = new ClientMessage.ReadAllEventsBackwardCompleted(Guid.NewGuid(), ReadAllResult.Success, "", - events, - new StreamMetadata(), true, 10, new TFPos(0, 0), - new TFPos(200, 200), new TFPos(0, 0), 100); - - var package = _dispatcher.WrapMessage(msg, (byte)ClientVersion.V2); - Assert.IsNotNull(package, "Package is null"); - Assert.AreEqual(TcpCommand.ReadAllEventsBackwardCompleted, package.Value.Command, "TcpCommand"); - - var dto = package.Value.Data.Deserialize(); - Assert.IsNotNull(dto, "DTO is null"); - Assert.AreEqual(1, dto.Events.Count(), "Number of events"); - - Assert.AreEqual(long.MaxValue, dto.Events[0].Event.EventNumber, "Event Number"); - } - - [Test] - public void - when_wrapping_read_all_events_backward_completed_with_link_to_deleted_event_should_not_downgrade_version_for_v2_clients() - { - var events = new ResolvedEvent[] { - ResolvedEvent.ForResolvedLink(CreateLinkEventRecord(), CreateDeletedEventRecord(), 100) - }; - var msg = new ClientMessage.ReadAllEventsBackwardCompleted(Guid.NewGuid(), ReadAllResult.Success, "", - events, - new StreamMetadata(), true, 10, new TFPos(0, 0), - new TFPos(200, 200), new TFPos(0, 0), 100); - - var package = _dispatcher.WrapMessage(msg, (byte)ClientVersion.V2); - Assert.IsNotNull(package, "Package is null"); - Assert.AreEqual(TcpCommand.ReadAllEventsBackwardCompleted, package.Value.Command, "TcpCommand"); - - var dto = package.Value.Data.Deserialize(); - Assert.IsNotNull(dto, "DTO is null"); - Assert.AreEqual(1, dto.Events.Count(), "Number of events"); - - Assert.AreEqual(0, dto.Events[0].Event.EventNumber, "Event Number"); - Assert.AreEqual(long.MaxValue, dto.Events[0].Link.EventNumber, "Link Event Number"); - } - - [Test] - public void - when_wrapping_stream_event_appeared_with_deleted_event_should_not_downgrade_version_for_v2_clients() - { - var msg = new ClientMessage.StreamEventAppeared(Guid.NewGuid(), - ResolvedEvent.ForUnresolvedEvent(CreateDeletedEventRecord(), 0)); - - var package = _dispatcher.WrapMessage(msg, (byte)ClientVersion.V2); - Assert.IsNotNull(package, "Package is null"); - Assert.AreEqual(TcpCommand.StreamEventAppeared, package.Value.Command, "TcpCommand"); - - var dto = package.Value.Data.Deserialize(); - Assert.IsNotNull(dto, "DTO is null"); - Assert.AreEqual(long.MaxValue, dto.Event.Event.EventNumber, "Event Number"); - } - - [Test] - public void - when_wrapping_subscribe_to_stream_confirmation_when_stream_deleted_should_not_downgrade_version_for_v2_clients() - { - var msg = new ClientMessage.SubscriptionConfirmation(Guid.NewGuid(), 100, long.MaxValue); - var package = _dispatcher.WrapMessage(msg, (byte)ClientVersion.V2); - Assert.IsNotNull(package, "Package is null"); - Assert.AreEqual(TcpCommand.SubscriptionConfirmation, package.Value.Command, "TcpCommand"); - - var dto = package.Value.Data.Deserialize(); - Assert.IsNotNull(dto, "DTO is null"); - Assert.AreEqual(long.MaxValue, dto.LastEventNumber, "Last Event Number"); - } - - [Test] - public void - when_wrapping_subscribe_to_stream_confirmation_when_stream_deleted_should_not_downgrade_last_event_number_for_v2_clients() - { - var msg = new ClientMessage.SubscriptionConfirmation(Guid.NewGuid(), 100, long.MaxValue); - var package = _dispatcher.WrapMessage(msg, (byte)ClientVersion.V2); - Assert.IsNotNull(package, "Package is null"); - Assert.AreEqual(TcpCommand.SubscriptionConfirmation, package.Value.Command, "TcpCommand"); - - var dto = package.Value.Data.Deserialize(); - Assert.IsNotNull(dto, "DTO is null"); - Assert.AreEqual(long.MaxValue, dto.LastEventNumber, "Last Event Number"); - } - - [Test] - public void - when_wrapping_stream_event_appeared_with_link_to_deleted_event_should_not_downgrade_version_for_v2_clients() - { - var msg = new ClientMessage.StreamEventAppeared(Guid.NewGuid(), - ResolvedEvent.ForResolvedLink(CreateLinkEventRecord(), CreateDeletedEventRecord(), 0)); - - var package = _dispatcher.WrapMessage(msg, (byte)ClientVersion.V2); - Assert.IsNotNull(package, "Package is null"); - Assert.AreEqual(TcpCommand.StreamEventAppeared, package.Value.Command, "TcpCommand"); - - var dto = package.Value.Data.Deserialize(); - Assert.IsNotNull(dto, "DTO is null"); - Assert.AreEqual(0, dto.Event.Event.EventNumber, "Event Number"); - Assert.AreEqual(long.MaxValue, dto.Event.Link.EventNumber, "Link Event Number"); - } - - [Test] - public void - when_wrapping_persistent_subscription_confirmation_when_stream_deleted_should_not_downgrade_last_event_number_for_v2_clients() - { - var msg = new ClientMessage.PersistentSubscriptionConfirmation("subscription", Guid.NewGuid(), 100, - long.MaxValue); - var package = _dispatcher.WrapMessage(msg, (byte)ClientVersion.V2); - Assert.IsNotNull(package, "Package is null"); - Assert.AreEqual(TcpCommand.PersistentSubscriptionConfirmation, package.Value.Command, "TcpCommand"); - - var dto = package.Value.Data.Deserialize(); - Assert.IsNotNull(dto, "DTO is null"); - Assert.AreEqual(long.MaxValue, dto.LastEventNumber, "Last event number"); - } - - [Test] - public void - when_wrapping_scavenge_started_response_should_return_result_and_scavengeId_for_v2_clients() - { - var scavengeId = Guid.NewGuid().ToString(); - var msg = new ClientMessage.ScavengeDatabaseStartedResponse(Guid.NewGuid(), scavengeId); - - var package = _dispatcher.WrapMessage(msg, (byte)ClientVersion.V2); - Assert.IsNotNull(package, "Package is null"); - Assert.AreEqual(TcpCommand.ScavengeDatabaseResponse, package.Value.Command, "TcpCommand"); - - var dto = package.Value.Data.Deserialize(); - Assert.IsNotNull(dto, "DTO is null"); - Assert.AreEqual(dto.Result, ScavengeDatabaseResponse.Types.ScavengeResult.Started); - Assert.AreEqual(dto.ScavengeId, scavengeId); - } - - [Test] - public void - when_wrapping_scavenge_inprogress_response_should_return_result_and_scavengeId_for_v2_clients() - { - var scavengeId = Guid.NewGuid().ToString(); - var msg = new ClientMessage.ScavengeDatabaseInProgressResponse(Guid.NewGuid(), scavengeId, reason: "In Progress"); - - var package = _dispatcher.WrapMessage(msg, (byte)ClientVersion.V2); - Assert.IsNotNull(package, "Package is null"); - Assert.AreEqual(TcpCommand.ScavengeDatabaseResponse, package.Value.Command, "TcpCommand"); - - var dto = package.Value.Data.Deserialize(); - Assert.IsNotNull(dto, "DTO is null"); - Assert.AreEqual(dto.Result, ScavengeDatabaseResponse.Types.ScavengeResult.InProgress); - Assert.AreEqual(dto.ScavengeId, scavengeId); - } - - [Test] - public void - when_wrapping_scavenge_unauthorized_response_should_return_result_and_scavengeId_for_v2_clients() - { - var scavengeId = Guid.NewGuid().ToString(); - var msg = new ClientMessage.ScavengeDatabaseUnauthorizedResponse(Guid.NewGuid(), scavengeId, "Unauthorized"); - - var package = _dispatcher.WrapMessage(msg, (byte)ClientVersion.V2); - Assert.IsNotNull(package, "Package is null"); - Assert.AreEqual(TcpCommand.ScavengeDatabaseResponse, package.Value.Command, "TcpCommand"); - - var dto = package.Value.Data.Deserialize(); - Assert.IsNotNull(dto, "DTO is null"); - Assert.AreEqual(dto.Result, ScavengeDatabaseResponse.Types.ScavengeResult.Unauthorized); - Assert.AreEqual(dto.ScavengeId, scavengeId); - } - - private EventRecord CreateDeletedEventRecord() - { - return new EventRecord(long.MaxValue, - LogRecord.DeleteTombstone(new LogV2RecordFactory(), 0, Guid.NewGuid(), Guid.NewGuid(), - "test-stream", "test-type", long.MaxValue), "test-stream", SystemEventTypes.StreamDeleted); - } - - private EventRecord CreateLinkEventRecord() - { - return new EventRecord(0, LogRecord.Prepare(new LogV2RecordFactory(), 100, Guid.NewGuid(), Guid.NewGuid(), 0, 0, - "link-stream", -1, PrepareFlags.SingleWrite | PrepareFlags.Data, SystemEventTypes.LinkTo, - Encoding.UTF8.GetBytes(string.Format("{0}@test-stream", long.MaxValue)), new byte[0]), "link-stream", SystemEventTypes.LinkTo); - } -} diff --git a/src/EventStore.Core.Tests/Services/Transport/Tcp/TcpConnectionManagerTests.cs b/src/EventStore.Core.Tests/Services/Transport/Tcp/TcpConnectionManagerTests.cs deleted file mode 100644 index ff724c7917..0000000000 --- a/src/EventStore.Core.Tests/Services/Transport/Tcp/TcpConnectionManagerTests.cs +++ /dev/null @@ -1,348 +0,0 @@ -using System; -using System.Collections.Generic; -using System.Linq; -using System.Net; -using System.Net.Sockets; -using System.Threading; -using EventStore.Client.Messages; -using EventStore.Core.Authentication.InternalAuthentication; -using EventStore.Core.Bus; -using EventStore.Core.Data; -using EventStore.Core.Messages; -using EventStore.Core.Messaging; -using EventStore.Core.Services; -using EventStore.Core.Services.Transport.Tcp; -using EventStore.Core.Settings; -using EventStore.Core.Tests.Authentication; -using EventStore.Core.Tests.Authorization; -using EventStore.Core.TransactionLog.LogRecords; -using EventStore.Core.Util; -using EventStore.Transport.Tcp; -using NUnit.Framework; -using EventRecord = EventStore.Core.Data.EventRecord; -using ResolvedEvent = EventStore.Core.Data.ResolvedEvent; - -namespace EventStore.Core.Tests.Services.Transport.Tcp; - -[TestFixture] -public class TcpConnectionManagerTests -{ - private int _connectionPendingSendBytesThreshold = 10 * 1024; - private int _connectionQueueSizeThreshold = 50000; - - [Test] - public void when_handling_trusted_write_on_external_service() - { - var package = new TcpPackage(TcpCommand.WriteEvents, TcpFlags.TrustedWrite, Guid.NewGuid(), null, null, - new byte[] { }); - - var dummyConnection = new DummyTcpConnection(); - - var tcpConnectionManager = new TcpConnectionManager( - Guid.NewGuid().ToString(), TcpServiceType.External, new ClientTcpDispatcher(2000), - new SynchronousScheduler(), dummyConnection, new SynchronousScheduler(), - new InternalAuthenticationProvider( - InMemoryBus.CreateTest(), new Core.Helpers.IODispatcher(new SynchronousScheduler(), new NoopEnvelope()), - new StubPasswordHashAlgorithm(), 1, false, DefaultData.DefaultUserOptions), - new AuthorizationGateway(new TestAuthorizationProvider()), - TimeSpan.FromSeconds(10), TimeSpan.FromSeconds(10), (man, err) => { }, - _connectionPendingSendBytesThreshold, _connectionQueueSizeThreshold); - - tcpConnectionManager.ProcessPackage(package); - - var data = dummyConnection.ReceivedData.Last(); - var receivedPackage = TcpPackage.FromArraySegment(data); - - Assert.AreEqual(receivedPackage.Command, TcpCommand.BadRequest, "Expected Bad Request but got {0}", - receivedPackage.Command); - } - - [Test] - public void when_handling_trusted_write_on_internal_service() - { - ManualResetEvent waiter = new ManualResetEvent(false); - ClientMessage.WriteEvents publishedWrite = null; - var evnt = new Event(Guid.NewGuid(), "TestEventType", true, new byte[] { }, new byte[] { }); - var write = new WriteEvents( - Guid.NewGuid().ToString(), - ExpectedVersion.Any, - new[] { - new NewEvent(evnt.EventId.ToByteArray(), evnt.EventType, evnt.IsJson ? 1 : 0, 0, - evnt.Data, evnt.Metadata) - }, - false); - - var package = new TcpPackage(TcpCommand.WriteEvents, Guid.NewGuid(), write.Serialize()); - var dummyConnection = new DummyTcpConnection(); - var publisher = new SynchronousScheduler(); - - publisher.Subscribe(new AdHocHandler(x => - { - publishedWrite = x; - waiter.Set(); - })); - - var tcpConnectionManager = new TcpConnectionManager( - Guid.NewGuid().ToString(), TcpServiceType.Internal, new ClientTcpDispatcher(2000), - publisher, dummyConnection, publisher, - new InternalAuthenticationProvider(publisher, new Core.Helpers.IODispatcher(publisher, new NoopEnvelope()), - new StubPasswordHashAlgorithm(), 1, false, DefaultData.DefaultUserOptions), - new AuthorizationGateway(new TestAuthorizationProvider()), - TimeSpan.FromSeconds(10), TimeSpan.FromSeconds(10), (man, err) => { }, - _connectionPendingSendBytesThreshold, _connectionQueueSizeThreshold); - - tcpConnectionManager.ProcessPackage(package); - - if (!waiter.WaitOne(TimeSpan.FromSeconds(5))) - { - throw new Exception("Timed out waiting for events."); - } - - Assert.AreEqual(evnt.EventId, publishedWrite.Events.First().EventId, - "Expected the published write to be the event that was sent through the tcp connection manager to be the event {0} but got {1}", - evnt.EventId, publishedWrite.Events.First().EventId); - } - - [Test] - public void - when_limit_pending_and_sending_message_smaller_than_threshold_and_pending_bytes_over_threshold_should_close_connection() - { - var mre = new ManualResetEventSlim(); - - var messageSize = _connectionPendingSendBytesThreshold / 2; - var evnt = new EventRecord(0, 0, Guid.NewGuid(), Guid.NewGuid(), 0, 0, "testStream", 0, DateTime.Now, - PrepareFlags.None, "eventType", new byte[messageSize], new byte[0]); - var record = ResolvedEvent.ForUnresolvedEvent(evnt, null); - var message = new ClientMessage.ReadEventCompleted(Guid.NewGuid(), "testStream", ReadEventResult.Success, - record, StreamMetadata.Empty, false, ""); - - var dummyConnection = new DummyTcpConnection(); - dummyConnection.PendingSendBytes = _connectionPendingSendBytesThreshold + 1000; - - var tcpConnectionManager = new TcpConnectionManager( - Guid.NewGuid().ToString(), TcpServiceType.External, new ClientTcpDispatcher(2000), - new SynchronousScheduler(), dummyConnection, new SynchronousScheduler(), - new InternalAuthenticationProvider( - InMemoryBus.CreateTest(), new Core.Helpers.IODispatcher(new SynchronousScheduler(), - new NoopEnvelope()), null, 1, false, DefaultData.DefaultUserOptions), - new AuthorizationGateway(new TestAuthorizationProvider()), - TimeSpan.FromSeconds(10), TimeSpan.FromSeconds(10), (man, err) => { mre.Set(); }, - _connectionPendingSendBytesThreshold, _connectionQueueSizeThreshold); - - tcpConnectionManager.SendMessage(message); - - if (!mre.Wait(2000)) - { - Assert.Fail("Timed out waiting for connection to close"); - } - } - - [Test] - public void - when_limit_pending_and_sending_message_larger_than_pending_bytes_threshold_but_no_bytes_pending_should_not_close_connection() - { - var messageSize = _connectionPendingSendBytesThreshold + 1000; - var evnt = new EventRecord(0, 0, Guid.NewGuid(), Guid.NewGuid(), 0, 0, "testStream", 0, DateTime.Now, - PrepareFlags.None, "eventType", new byte[messageSize], new byte[0]); - var record = ResolvedEvent.ForUnresolvedEvent(evnt, null); - var message = new ClientMessage.ReadEventCompleted(Guid.NewGuid(), "testStream", ReadEventResult.Success, - record, StreamMetadata.Empty, false, ""); - - var dummyConnection = new DummyTcpConnection(); - dummyConnection.PendingSendBytes = 0; - - var tcpConnectionManager = new TcpConnectionManager( - Guid.NewGuid().ToString(), TcpServiceType.External, new ClientTcpDispatcher(2000), - new SynchronousScheduler(), dummyConnection, new SynchronousScheduler(), - new InternalAuthenticationProvider(InMemoryBus.CreateTest(), - new Core.Helpers.IODispatcher(new SynchronousScheduler(), new NoopEnvelope()), null, 1, false, DefaultData.DefaultUserOptions), - new AuthorizationGateway(new TestAuthorizationProvider()), - TimeSpan.FromSeconds(10), TimeSpan.FromSeconds(10), (man, err) => { }, - _connectionPendingSendBytesThreshold, _connectionQueueSizeThreshold); - - tcpConnectionManager.SendMessage(message); - - var data = dummyConnection.ReceivedData.Last(); - var receivedPackage = TcpPackage.FromArraySegment(data); - - Assert.AreEqual(receivedPackage.Command, TcpCommand.ReadEventCompleted, - "Expected ReadEventCompleted but got {0}", receivedPackage.Command); - } - - [Test] - public void - when_not_limit_pending_and_sending_message_smaller_than_threshold_and_pending_bytes_over_threshold_should_not_close_connection() - { - var mre = new ManualResetEventSlim(); - - var messageSize = _connectionPendingSendBytesThreshold / 2; - var evnt = new EventRecord(0, 0, Guid.NewGuid(), Guid.NewGuid(), 0, 0, "testStream", 0, DateTime.Now, - PrepareFlags.None, "eventType", new byte[messageSize], new byte[0]); - var record = ResolvedEvent.ForUnresolvedEvent(evnt, null); - var message = new ClientMessage.ReadEventCompleted(Guid.NewGuid(), "testStream", ReadEventResult.Success, - record, StreamMetadata.Empty, false, ""); - - var dummyConnection = new DummyTcpConnection(); - dummyConnection.PendingSendBytes = _connectionPendingSendBytesThreshold + 1000; - - var tcpConnectionManager = new TcpConnectionManager( - Guid.NewGuid().ToString(), TcpServiceType.External, new ClientTcpDispatcher(2000), - new SynchronousScheduler(), dummyConnection, new SynchronousScheduler(), - new InternalAuthenticationProvider(InMemoryBus.CreateTest(), - new Core.Helpers.IODispatcher(new SynchronousScheduler(), new NoopEnvelope()), null, 1, false, DefaultData.DefaultUserOptions), - new AuthorizationGateway(new TestAuthorizationProvider()), - TimeSpan.FromSeconds(10), TimeSpan.FromSeconds(10), (man, err) => { mre.Set(); }, - ESConsts.UnrestrictedPendingSendBytes, ESConsts.MaxConnectionQueueSize); - - tcpConnectionManager.SendMessage(message); - - var data = dummyConnection.ReceivedData.Last(); - var receivedPackage = TcpPackage.FromArraySegment(data); - - Assert.AreEqual(receivedPackage.Command, TcpCommand.ReadEventCompleted, - "Expected ReadEventCompleted but got {0}", receivedPackage.Command); - } - - [Test] - public void - when_send_queue_size_is_smaller_than_threshold_should_not_close_connection() - { - var mre = new ManualResetEventSlim(); - - var messageSize = _connectionPendingSendBytesThreshold; - var evnt = new EventRecord(0, 0, Guid.NewGuid(), Guid.NewGuid(), 0, 0, "testStream", 0, DateTime.Now, - PrepareFlags.None, "eventType", new byte[messageSize], new byte[0]); - var record = ResolvedEvent.ForUnresolvedEvent(evnt, null); - var message = new ClientMessage.ReadEventCompleted(Guid.NewGuid(), "testStream", ReadEventResult.Success, - record, StreamMetadata.Empty, false, ""); - - var dummyConnection = new DummyTcpConnection(); - dummyConnection.SendQueueSize = ESConsts.MaxConnectionQueueSize - 1; - - var tcpConnectionManager = new TcpConnectionManager( - Guid.NewGuid().ToString(), TcpServiceType.External, new ClientTcpDispatcher(2000), - new SynchronousScheduler(), dummyConnection, new SynchronousScheduler(), - new InternalAuthenticationProvider(InMemoryBus.CreateTest(), - new Core.Helpers.IODispatcher(new SynchronousScheduler(), new NoopEnvelope()), null, 1, false, DefaultData.DefaultUserOptions), - new AuthorizationGateway(new TestAuthorizationProvider()), - TimeSpan.FromSeconds(10), TimeSpan.FromSeconds(10), (man, err) => { mre.Set(); }, - ESConsts.UnrestrictedPendingSendBytes, ESConsts.MaxConnectionQueueSize); - - tcpConnectionManager.SendMessage(message); - - var data = dummyConnection.ReceivedData.Last(); - var receivedPackage = TcpPackage.FromArraySegment(data); - - Assert.AreEqual(receivedPackage.Command, TcpCommand.ReadEventCompleted, - "Expected ReadEventCompleted but got {0}", receivedPackage.Command); - } - - [Test] - public void - when_send_queue_size_is_larger_than_threshold_should_close_connection() - { - var mre = new ManualResetEventSlim(); - - var messageSize = _connectionPendingSendBytesThreshold; - var evnt = new EventRecord(0, 0, Guid.NewGuid(), Guid.NewGuid(), 0, 0, "testStream", 0, DateTime.Now, - PrepareFlags.None, "eventType", new byte[messageSize], new byte[0]); - var record = ResolvedEvent.ForUnresolvedEvent(evnt, null); - var message = new ClientMessage.ReadEventCompleted(Guid.NewGuid(), "testStream", ReadEventResult.Success, - record, StreamMetadata.Empty, false, ""); - - var dummyConnection = new DummyTcpConnection(); - dummyConnection.SendQueueSize = ESConsts.MaxConnectionQueueSize + 1; - - var tcpConnectionManager = new TcpConnectionManager( - Guid.NewGuid().ToString(), TcpServiceType.External, new ClientTcpDispatcher(2000), - new SynchronousScheduler(), dummyConnection, new SynchronousScheduler(), - new InternalAuthenticationProvider(InMemoryBus.CreateTest(), - new Core.Helpers.IODispatcher(new SynchronousScheduler(), new NoopEnvelope()), null, 1, false, DefaultData.DefaultUserOptions), - new AuthorizationGateway(new TestAuthorizationProvider()), - TimeSpan.FromSeconds(10), TimeSpan.FromSeconds(10), (man, err) => { mre.Set(); }, - ESConsts.UnrestrictedPendingSendBytes, ESConsts.MaxConnectionQueueSize); - - tcpConnectionManager.SendMessage(message); - - if (!mre.Wait(2000)) - { - Assert.Fail("Timed out waiting for connection to close"); - } - } -} - -internal class DummyTcpConnection : ITcpConnection -{ - public Guid ConnectionId - { - get { return _connectionId; } - set { _connectionId = value; } - } - private Guid _connectionId = Guid.NewGuid(); - public string ClientConnectionName - { - get { return _clientConnectionName; } - } - - public long TotalBytesSent { get; } - public long TotalBytesReceived { get; } - - public bool IsClosed - { - get { return false; } - } - - public IPEndPoint LocalEndPoint - { - get { return new IPEndPoint(IPAddress.Loopback, 2); } - } - - public IPEndPoint RemoteEndPoint - { - get { return new IPEndPoint(IPAddress.Loopback, 1); } - } - - private int _sendQueueSize; - public int SendQueueSize - { - get { return _sendQueueSize; } - set { _sendQueueSize = value; } - } - - private int _pendingSendBytes; - - public int PendingSendBytes - { - get { return _pendingSendBytes; } - set { _pendingSendBytes = value; } - } - - public event Action ConnectionClosed; - private string _clientConnectionName; - - public void Close(string reason) - { - var handler = ConnectionClosed; - if (handler != null) - { - handler(this, SocketError.Shutdown); - } - } - - public IEnumerable> ReceivedData; - - public void EnqueueSend(IEnumerable> data) - { - ReceivedData = data; - } - - public void ReceiveAsync(Action>> callback) - { - throw new NotImplementedException(); - } - - public void SetClientConnectionName(string clientConnectionName) - { - _clientConnectionName = clientConnectionName; - } -} diff --git a/src/EventStore.Core.Tests/Services/Transport/Tcp/TcpConnectionSslTests.cs b/src/EventStore.Core.Tests/Services/Transport/Tcp/TcpConnectionSslTests.cs deleted file mode 100644 index 486d0462de..0000000000 --- a/src/EventStore.Core.Tests/Services/Transport/Tcp/TcpConnectionSslTests.cs +++ /dev/null @@ -1,176 +0,0 @@ -using System; -using System.Collections.Generic; -using System.IO; -using System.Net; -using System.Net.Sockets; -using System.Reflection; -using System.Security.Cryptography.X509Certificates; -using System.Threading; -using System.Threading.Tasks; -using EventStore.Common.Utils; -using EventStore.Transport.Tcp; -using NUnit.Framework; - -namespace EventStore.Core.Tests.Services.Transport.Tcp; - -[TestFixture] -public class TcpConnectionSslTests -{ - protected static Socket CreateListeningSocket() - { - var listener = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); - listener.Bind(new IPEndPoint(IPAddress.Loopback, 0)); - listener.Listen(1); - return listener; - } - - private IEnumerable> GenerateData() - { - var data = new List>(); - data.Add(new ArraySegment(new byte[100])); - return data; - } - - [Test, Timeout(120000)] - public async Task no_data_should_be_dispatched_after_tcp_connection_closed() - { - for (int i = 0; i < 1000; i++) - { - bool closed = false; - bool dataReceivedAfterClose = false; - var listeningSocket = CreateListeningSocket(); - - var mre = new ManualResetEventSlim(false); - var clientTcpConnection = TcpConnectionSsl.CreateConnectingConnection( - Guid.NewGuid(), - listeningSocket.LocalEndPoint.GetHost(), - null, - (IPEndPoint)listeningSocket.LocalEndPoint, - delegate - { return (true, null); }, - null, - new TcpClientConnector(), - TimeSpan.FromSeconds(5), - (conn) => mre.Set(), - (conn, error) => - { - Assert.Fail($"Connection failed: {error}"); - }, - false); - - var serverSocket = listeningSocket.Accept(); - var serverTcpConnection = TcpConnectionSsl.CreateServerFromSocket(Guid.NewGuid(), - (IPEndPoint)serverSocket.RemoteEndPoint, serverSocket, ssl_connections.GetServerCertificate, - null, delegate - { return (true, null); }, false); - - mre.Wait(TimeSpan.FromSeconds(3)); - try - { - clientTcpConnection.ConnectionClosed += (connection, error) => - { - Volatile.Write(ref closed, true); - }; - - clientTcpConnection.ReceiveAsync((connection, data) => - { - if (Volatile.Read(ref closed)) - { - dataReceivedAfterClose = true; - } - }); - - using (var b = new Barrier(2)) - { - Task sendData = Task.Factory.StartNew(() => - { - b.SignalAndWait(); - for (int i = 0; i < 1000; i++) - { - serverTcpConnection.EnqueueSend(GenerateData()); - } - }, CancellationToken.None, TaskCreationOptions.LongRunning, TaskScheduler.Default); - - Task closeConnection = Task.Factory.StartNew(() => - { - b.SignalAndWait(); - serverTcpConnection.Close("Intentional close"); - }, CancellationToken.None, TaskCreationOptions.LongRunning, TaskScheduler.Default); - - await Task.WhenAll(sendData, closeConnection); - Assert.False(dataReceivedAfterClose); - } - } - finally - { - clientTcpConnection.Close("Shut down"); - serverTcpConnection.Close("Shut down"); - listeningSocket.Dispose(); - } - } - } - - [Test, Timeout(120000)] - public void when_connection_closed_quickly_socket_should_be_properly_disposed() - { - for (int i = 0; i < 1000; i++) - { - var listeningSocket = CreateListeningSocket(); - ITcpConnection clientTcpConnection = null; - ITcpConnection serverTcpConnection = null; - Socket serverSocket = null; - try - { - ManualResetEventSlim mre = new ManualResetEventSlim(false); - - clientTcpConnection = TcpConnectionSsl.CreateConnectingConnection( - Guid.NewGuid(), - listeningSocket.LocalEndPoint.GetHost(), - null, - (IPEndPoint)listeningSocket.LocalEndPoint, - delegate - { return (true, null); }, - null, - new TcpClientConnector(), - TimeSpan.FromSeconds(5), - (conn) => { }, - (conn, error) => { }, - false); - - clientTcpConnection.ConnectionClosed += (conn, error) => - { - mre.Set(); - }; - - serverSocket = listeningSocket.Accept(); - clientTcpConnection.Close("Intentional close"); - serverTcpConnection = TcpConnectionSsl.CreateServerFromSocket(Guid.NewGuid(), - (IPEndPoint)serverSocket.RemoteEndPoint, serverSocket, ssl_connections.GetServerCertificate, - null, delegate - { return (true, null); }, false); - - mre.Wait(TimeSpan.FromSeconds(10)); - SpinWait.SpinUntil(() => serverTcpConnection.IsClosed, TimeSpan.FromSeconds(10)); - - var disposed = false; - try - { - int x = serverSocket.Available; - } - catch (ObjectDisposedException) - { - disposed = true; - } - - Assert.AreEqual(true, disposed); - } - finally - { - clientTcpConnection?.Close("Shut down"); - serverTcpConnection?.Close("Shut down"); - listeningSocket.Dispose(); - serverSocket?.Dispose(); - } - } - } -} diff --git a/src/EventStore.Core.Tests/Services/Transport/Tcp/TcpConnectionTests.cs b/src/EventStore.Core.Tests/Services/Transport/Tcp/TcpConnectionTests.cs deleted file mode 100644 index fb40e56fdd..0000000000 --- a/src/EventStore.Core.Tests/Services/Transport/Tcp/TcpConnectionTests.cs +++ /dev/null @@ -1,156 +0,0 @@ -using System; -using System.Collections.Generic; -using System.Net; -using System.Net.Sockets; -using System.Threading; -using System.Threading.Tasks; -using EventStore.Transport.Tcp; -using NUnit.Framework; - -namespace EventStore.Core.Tests.Services.Transport.Tcp; - -[TestFixture] -public class TcpConnectionTests -{ - protected static Socket CreateListeningSocket() - { - var listener = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp); - listener.Bind(new IPEndPoint(IPAddress.Loopback, 0)); - listener.Listen(1); - return listener; - } - - private IEnumerable> GenerateData() - { - var data = new List>(); - data.Add(new ArraySegment(new byte[100])); - return data; - } - - [Test, Timeout(120000)] - public async Task no_data_should_be_dispatched_after_tcp_connection_closed() - { - for (int i = 0; i < 1000; i++) - { - bool closed = false; - bool dataReceivedAfterClose = false; - var listeningSocket = CreateListeningSocket(); - TaskCompletionSource connectionResult = new(TaskCreationOptions.RunContinuationsAsynchronously); - - var clientTcpConnection = TcpConnection.CreateConnectingTcpConnection( - Guid.NewGuid(), - (IPEndPoint)listeningSocket.LocalEndPoint, - new TcpClientConnector(), - TimeSpan.FromSeconds(5), - (conn) => connectionResult.TrySetResult(SocketError.Success), - (conn, error) => connectionResult.TrySetResult(error), - false); - - var serverSocket = listeningSocket.Accept(); - var serverTcpConnection = TcpConnection.CreateAcceptedTcpConnection(Guid.NewGuid(), - (IPEndPoint)serverSocket.RemoteEndPoint, serverSocket, false); - - SocketError error = await connectionResult.Task.WithTimeout(); - Assert.AreEqual(SocketError.Success, error); - try - { - clientTcpConnection.ConnectionClosed += (connection, error) => - { - Volatile.Write(ref closed, true); - }; - - clientTcpConnection.ReceiveAsync((connection, data) => - { - if (Volatile.Read(ref closed)) - { - dataReceivedAfterClose = true; - } - }); - - using (var b = new Barrier(2)) - { - Task sendData = Task.Factory.StartNew(() => - { - b.SignalAndWait(); - for (int i = 0; i < 1000; i++) - { - serverTcpConnection.EnqueueSend(GenerateData()); - } - }, CancellationToken.None, TaskCreationOptions.LongRunning, TaskScheduler.Default); - - Task closeConnection = Task.Factory.StartNew(() => - { - b.SignalAndWait(); - serverTcpConnection.Close("Intentional close"); - }, CancellationToken.None, TaskCreationOptions.LongRunning, TaskScheduler.Default); - - await Task.WhenAll(sendData, closeConnection); - Assert.False(dataReceivedAfterClose); - } - } - finally - { - clientTcpConnection.Close("Shut down"); - serverTcpConnection.Close("Shut down"); - listeningSocket.Dispose(); - } - } - } - - [Test, Timeout(120000)] - public void when_connection_closed_quickly_socket_should_be_properly_disposed() - { - for (int i = 0; i < 1000; i++) - { - var listeningSocket = CreateListeningSocket(); - ITcpConnection clientTcpConnection = null; - ITcpConnection serverTcpConnection = null; - Socket serverSocket = null; - try - { - ManualResetEventSlim mre = new ManualResetEventSlim(false); - - clientTcpConnection = TcpConnection.CreateConnectingTcpConnection( - Guid.NewGuid(), - (IPEndPoint)listeningSocket.LocalEndPoint, - new TcpClientConnector(), - TimeSpan.FromSeconds(5), - (conn) => { }, - (conn, error) => { }, - false); - - clientTcpConnection.ConnectionClosed += (conn, error) => - { - mre.Set(); - }; - - serverSocket = listeningSocket.Accept(); - clientTcpConnection.Close("Intentional close"); - serverTcpConnection = TcpConnection.CreateAcceptedTcpConnection(Guid.NewGuid(), - (IPEndPoint)serverSocket.RemoteEndPoint, serverSocket, false); - - mre.Wait(TimeSpan.FromSeconds(10)); - SpinWait.SpinUntil(() => serverTcpConnection.IsClosed, TimeSpan.FromSeconds(10)); - - var disposed = false; - try - { - int x = serverSocket.Available; - } - catch (ObjectDisposedException) - { - disposed = true; - } - - Assert.AreEqual(true, disposed); - } - finally - { - clientTcpConnection?.Close("Shut down"); - serverTcpConnection?.Close("Shut down"); - listeningSocket.Dispose(); - serverSocket?.Dispose(); - } - } - } -} diff --git a/src/EventStore.Core.Tests/Services/Transport/Tcp/when_invalid_data_is_sent_over_tcp.cs b/src/EventStore.Core.Tests/Services/Transport/Tcp/when_invalid_data_is_sent_over_tcp.cs deleted file mode 100644 index 66ca036d1f..0000000000 --- a/src/EventStore.Core.Tests/Services/Transport/Tcp/when_invalid_data_is_sent_over_tcp.cs +++ /dev/null @@ -1,90 +0,0 @@ -using System; -using System.Collections.Generic; -using System.Net; -using System.Net.Sockets; -using System.Threading; -using System.Threading.Tasks; -using EventStore.Common.Utils; -using EventStore.Core.Tests.Integration; -using EventStore.Transport.Tcp; -using NUnit.Framework; - -namespace EventStore.Core.Tests.Services.Transport.Tcp; - -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class when_invalid_data_is_sent_over_tcp : specification_with_cluster -{ - - [Timeout(15000)] - [TestCase("ExternalTcpEndPoint", false)] - public async Task connection_should_be_closed_by_remote_party(string endpointProperty, bool secure) - { - IPEndPoint endpoint = (IPEndPoint)_nodes[0].GetType().GetProperty(endpointProperty).GetValue(_nodes[0], null); - await WaitForEndpoint(endpoint); - - var closedEvent = new ManualResetEventSlim(); - TaskCompletionSource connectionResult = new(TaskCreationOptions.RunContinuationsAsynchronously); - - ITcpConnection connection; - if (!secure) - { - connection = TcpConnection.CreateConnectingTcpConnection( - Guid.NewGuid(), - endpoint, - new TcpClientConnector(), - TimeSpan.FromSeconds(5), - (conn) => connectionResult.TrySetResult(SocketError.Success), - (conn, error) => connectionResult.TrySetResult(error), - false); - } - else - { - connection = TcpConnectionSsl.CreateConnectingConnection( - Guid.NewGuid(), - endpoint.GetHost(), - null, - endpoint, - delegate - { return (true, null); }, - null, - new TcpClientConnector(), - TimeSpan.FromSeconds(5), - (conn) => connectionResult.TrySetResult(SocketError.Success), - (conn, error) => connectionResult.TrySetResult(error), - false); - } - - connection.ConnectionClosed += (conn, error) => closedEvent.Set(); - - SocketError result = await connectionResult.Task.WithTimeout(); - Assert.AreEqual(SocketError.Success, result); - var data = new List> { - new ArraySegment(new byte[] {1, 2, 3, 4, 5, 6, 7, 8, 9, 10}) - }; - connection.EnqueueSend(data); - Assert.True(closedEvent.Wait(10000)); - connection.Close("intentional close"); - } - - private static async Task WaitForEndpoint(IPEndPoint endpoint) - { - using var timeout = new CancellationTokenSource(TimeSpan.FromSeconds(5)); - - while (!timeout.IsCancellationRequested) - { - using var client = new TcpClient(); - - try - { - await client.ConnectAsync(endpoint.Address, endpoint.Port, timeout.Token); - return; - } - catch (Exception ex) when (ex is SocketException or OperationCanceledException) - { - await Task.Delay(100, CancellationToken.None); - } - } - - throw new TimeoutException($"TCP endpoint {endpoint} did not accept connections before the test timeout."); - } -} diff --git a/src/EventStore.Core.Tests/Services/Transport/Tcp/with_intermediate_certificates.cs b/src/EventStore.Core.Tests/Services/Transport/Tcp/with_intermediate_certificates.cs deleted file mode 100644 index 2b1c2f1540..0000000000 --- a/src/EventStore.Core.Tests/Services/Transport/Tcp/with_intermediate_certificates.cs +++ /dev/null @@ -1,156 +0,0 @@ -using System; -using System.Net; -using System.Net.Security; -using System.Security.Cryptography.X509Certificates; -using System.Threading; -using EventStore.Common.Utils; -using EventStore.Core.Services.Transport.Tcp; -using EventStore.Core.Tests.Certificates; -using EventStore.Core.Tests.Helpers; -using EventStore.Transport.Tcp; -using NUnit.Framework; - -namespace EventStore.Core.Tests.Services.Transport.Tcp; - -[TestFixture] -public class with_intermediate_certificates : with_certificate_chain_of_length_3 -{ - private TcpServerListener _listener; - private IPEndPoint _serverEndPoint; - private ITcpConnection _client; - private Func _clientCertValidator; - private X509Certificate2 _cert; - - [SetUp] - public void SetUp() - { - // certificate exported to PKCS #12 due to this issue on Windows: https://github.com/dotnet/runtime/issues/45680 - _cert = X509CertificateLoader.LoadPkcs12(_leaf.ExportToPkcs12(), null); - - _clientCertValidator = (_, _, _) => (true, null); - _serverEndPoint = new IPEndPoint(IPAddress.Loopback, PortsHelper.GetAvailablePort(IPAddress.Loopback)); - _listener = new TcpServerListener(_serverEndPoint); - _listener.StartListening((endPoint, socket) => - { - TcpConnectionSsl.CreateServerFromSocket( - Guid.NewGuid(), - endPoint, - socket, - () => _cert, - () => new X509Certificate2Collection(_intermediate), - (cert, chain, errors) => _clientCertValidator(cert, chain, errors), - verbose: true); - }, "Secure"); - } - - [Test] - public void server_should_send_intermediate_certificate_during_handshake() - { - var done = new ManualResetEventSlim(false); - - bool gotLeaf = false; - bool gotIntermediate = false; - - _client = TcpConnectionSsl.CreateConnectingConnection( - Guid.NewGuid(), - _serverEndPoint.GetHost(), - null, - _serverEndPoint, - (certificate, chain, _, _) => - { - gotLeaf = _leaf.Equals(certificate); - foreach (var chainElement in chain.ChainElements) - { - if (chainElement.Certificate.Equals(_intermediate)) - { - gotIntermediate = true; - } - } - - done.Set(); - return (true, null); - }, - null, - new TcpClientConnector(), - TcpConnectionManager.ConnectionTimeout, - conn => { }, - (conn, err) => { }, - verbose: true); - - Assert.True(done.Wait(20000), "Took too long to receive completion."); - Assert.True(gotLeaf); - Assert.True(gotIntermediate); - } - - [Test, Ignore("Skipped since it adds an intermediate certificate to the current user's store")] - public void client_should_send_intermediate_certificate_during_handshake() - { - try - { - // see: https://github.com/dotnet/runtime/issues/47680#issuecomment-771093045 - AddIntermediateCertificateToStore(); - - var done = new ManualResetEventSlim(false); - - bool gotLeaf = false; - bool gotIntermediate = false; - - _clientCertValidator = (certificate, chain, _) => - { - gotLeaf = _leaf.Equals(certificate); - foreach (var chainElement in chain.ChainElements) - { - if (chainElement.Certificate.Equals(_intermediate)) - { - gotIntermediate = true; - } - } - - done.Set(); - return (true, null); - }; - - _client = TcpConnectionSsl.CreateConnectingConnection( - Guid.NewGuid(), - _serverEndPoint.GetHost(), - null, - _serverEndPoint, - (_, _, _, _) => (true, null), - () => new X509Certificate2Collection(_cert), - new TcpClientConnector(), - TcpConnectionManager.ConnectionTimeout, - conn => { }, - (conn, err) => { }, - verbose: true); - - Assert.True(done.Wait(20000), "Took too long to receive completion."); - Assert.True(gotLeaf); - Assert.True(gotIntermediate); - } - finally - { - RemoveIntermediateCertificateFromStore(); - } - } - - private void AddIntermediateCertificateToStore() - { - using var intermediateStore = new X509Store(StoreName.CertificateAuthority, StoreLocation.CurrentUser); - intermediateStore.Open(OpenFlags.ReadWrite); - intermediateStore.Add(_intermediate); - } - - private void RemoveIntermediateCertificateFromStore() - { - using var intermediateStore = new X509Store(StoreName.CertificateAuthority, StoreLocation.CurrentUser); - intermediateStore.Open(OpenFlags.ReadWrite); - intermediateStore.Remove(_intermediate); - } - - [TearDown] - public void TearDown() - { - _listener.Stop(); - _client.Close("Normal close."); - } -} diff --git a/src/EventStore.Core.Tests/copying_metadata.cs b/src/EventStore.Core.Tests/copying_metadata.cs deleted file mode 100644 index 502e451223..0000000000 --- a/src/EventStore.Core.Tests/copying_metadata.cs +++ /dev/null @@ -1,78 +0,0 @@ -using System; -using System.Collections.Generic; -using System.Linq; -using System.Text; -using EventStore.ClientAPI; -using NUnit.Framework; - -namespace EventStore.Core.Tests; - -[TestFixture] -public class copying_metadata -{ - [Test] - public void copies_empty_metadata() - { - var empty = StreamMetadata.Build().Build(); - var copied = empty.Copy().Build(); - Assert.AreEqual(empty.AsJsonString(), copied.AsJsonString()); - } - - [Test] - public void copies_all_values() - { - var source = StreamMetadata.Build() - .SetCacheControl(TimeSpan.FromDays(1)) - .SetCustomProperty("Test", "Value") - .SetReadRole("foo") - .SetWriteRole("bar") - .SetDeleteRole("baz") - .SetMetadataReadRole("qux") - .SetMetadataWriteRole("quux") - .SetMaxAge(TimeSpan.FromHours(1)) - .SetMaxCount(2) - .SetTruncateBefore(4) - .Build(); - var copied = source.Copy().Build(); - Assert.AreEqual(source.AsJsonString(), copied.AsJsonString()); - } - - [Test] - public void can_mutate_copy() - { - var source = StreamMetadata.Build() - .SetCacheControl(TimeSpan.FromDays(1)) - .SetCustomProperty("Test", "Value") - .SetReadRole("foo") - .SetWriteRole("bar") - .SetDeleteRole("baz") - .SetMetadataReadRole("qux") - .SetMetadataWriteRole("quux") - .SetMaxAge(TimeSpan.FromHours(1)) - .SetMaxCount(2) - .SetTruncateBefore(4) - .Build(); - - var expected = StreamMetadata.Build() - .SetCacheControl(TimeSpan.FromDays(1)) - .SetCustomProperty("Test", "Value") - .SetCustomProperty("Test2", "Value2") - .SetReadRole("foo") - .SetWriteRole("bar") - .SetDeleteRole("baz") - .SetMetadataReadRole("qux") - .SetMetadataWriteRole("quux") - .SetMaxAge(TimeSpan.FromHours(1)) - .SetMaxCount(4) - .SetTruncateBefore(4) - .Build(); - - - var copied = source.Copy() - .SetMaxCount(4) - .SetCustomProperty("Test2", "Value2") - .Build(); - - Assert.AreEqual(expected.AsJsonString(), copied.AsJsonString()); - } -}