From 2d332c1f4d6f1e49243f64c7e634d9326b681489 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Sat, 12 Sep 2026 18:52:38 -0400 Subject: [PATCH] chore(tests): retire obsolete TCP integration fixtures Signed-off-by: Yordis Prieto --- ...node_becomes_leader_with_unindexed_data.cs | 370 ------------------ .../Transport/Tcp/TcpClientDispatcherTests.cs | 317 --------------- .../Tcp/TcpConnectionManagerTests.cs | 348 ---------------- .../Transport/Tcp/TcpConnectionSslTests.cs | 176 --------- .../Transport/Tcp/TcpConnectionTests.cs | 156 -------- .../Tcp/when_invalid_data_is_sent_over_tcp.cs | 90 ----- .../Tcp/with_intermediate_certificates.cs | 156 -------- src/EventStore.Core.Tests/copying_metadata.cs | 78 ---- 8 files changed, 1691 deletions(-) delete mode 100644 src/EventStore.Core.Tests/Integration/when_node_becomes_leader_with_unindexed_data.cs delete mode 100644 src/EventStore.Core.Tests/Services/Transport/Tcp/TcpClientDispatcherTests.cs delete mode 100644 src/EventStore.Core.Tests/Services/Transport/Tcp/TcpConnectionManagerTests.cs delete mode 100644 src/EventStore.Core.Tests/Services/Transport/Tcp/TcpConnectionSslTests.cs delete mode 100644 src/EventStore.Core.Tests/Services/Transport/Tcp/TcpConnectionTests.cs delete mode 100644 src/EventStore.Core.Tests/Services/Transport/Tcp/when_invalid_data_is_sent_over_tcp.cs delete mode 100644 src/EventStore.Core.Tests/Services/Transport/Tcp/with_intermediate_certificates.cs delete mode 100644 src/EventStore.Core.Tests/copying_metadata.cs diff --git a/src/EventStore.Core.Tests/Integration/when_node_becomes_leader_with_unindexed_data.cs b/src/EventStore.Core.Tests/Integration/when_node_becomes_leader_with_unindexed_data.cs deleted file mode 100644 index f7e6400840..0000000000 --- a/src/EventStore.Core.Tests/Integration/when_node_becomes_leader_with_unindexed_data.cs +++ /dev/null @@ -1,370 +0,0 @@ -using System; -using System.Collections.Generic; -using System.Linq; -using System.Net; -using System.Net.Http; -using System.Text; -using System.Threading.Tasks; -using EventStore.Client.Streams; -using EventStore.Common.Utils; -using EventStore.Core.Data; -using EventStore.Core.Services.Transport.Grpc; -using EventStore.Core.Tests.Helpers; -using EventStore.Plugins.Subsystems; -using Google.Protobuf; -using Grpc.Core; -using Grpc.Net.Client; -using NUnit.Framework; -using Empty = EventStore.Client.Empty; -using GrpcMetadata = EventStore.Core.Services.Transport.Grpc.Constants.Metadata; -using RecordedEvent = EventStore.Client.Streams.ReadResp.Types.ReadEvent.Types.RecordedEvent; - -namespace EventStore.Core.Tests.Integration; - -[Explicit] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class when_node_becomes_leader_with_unindexed_data : specification_with_cluster -{ - private const string FakeHostAdvertiseAs = "192.168.123.123"; - private const string Username = "admin"; - private const string Password = "changeit"; - - private const string AuthenticationScheme = "Basic"; - private readonly string AuthenticationValue = Convert.ToBase64String(Encoding.ASCII.GetBytes($"{Username}:{Password}")); - private static readonly Marshaller EmptyMarshaller = Marshallers.Create( - static message => message.ToByteArray(), - static bytes => Empty.Parser.ParseFrom(bytes)); - private static readonly Method ResignNodeMethod = new( - MethodType.Unary, - "event_store.client.operations.Operations", - "ResignNode", - EmptyMarshaller, - EmptyMarshaller); - private class CommitTimeoutException : Exception { } - private class WrongExpectedVersionException : Exception { } - - private EndPoint[][] _nodeGossipSeeds; - private HttpClient _httpClient; - - protected override async Task Given() - { - _nodeGossipSeeds = new[] { - new EndPoint[] {_nodeEndpoints[1].HttpEndPoint, _nodeEndpoints[2].HttpEndPoint}, - new EndPoint[] {_nodeEndpoints[0].HttpEndPoint, _nodeEndpoints[2].HttpEndPoint}, - new EndPoint[] {_nodeEndpoints[0].HttpEndPoint, _nodeEndpoints[1].HttpEndPoint} - }; - _httpClient = new HttpClient(new SocketsHttpHandler - { - SslOptions = { - RemoteCertificateValidationCallback = delegate { return true; } - } - }, true); - await base.Given(); - } - - private CallOptions GetCallOptions() - { - return new( - credentials: CallCredentials.FromInterceptor((_, metadata) => - { - metadata.Add("authorization", $"{AuthenticationScheme} {AuthenticationValue}"); - return Task.CompletedTask; - }), - deadline: DateTime.UtcNow.AddSeconds(5)); - } - - private async Task AppendEvent(IPEndPoint endpoint, string stream, long expectedVersion) - { - using var channel = GrpcChannel.ForAddress(new Uri($"https://{endpoint}"), - new GrpcChannelOptions { HttpClient = _httpClient }); - var streamClient = new Streams.StreamsClient(channel); - using var call = streamClient.Append(GetCallOptions()); - - var optionsAppendReq = new AppendReq - { - Options = new() - { - StreamIdentifier = new() - { - StreamName = ByteString.CopyFromUtf8(stream) - }, - } - }; - switch (expectedVersion) - { - case ExpectedVersion.Any: - optionsAppendReq.Options.Any = new Empty(); - break; - case ExpectedVersion.NoStream: - optionsAppendReq.Options.NoStream = new Empty(); - break; - default: - optionsAppendReq.Options.Revision = (ulong)expectedVersion; - break; - } - - await call.RequestStream.WriteAsync(optionsAppendReq); - await call.RequestStream.WriteAsync(new AppendReq - { - ProposedMessage = new() - { - Id = new() - { - String = Uuid.FromGuid(Guid.NewGuid()).ToString() - }, - Data = ByteString.Empty, - Metadata = { - [GrpcMetadata.Type] = "type", - [GrpcMetadata.ContentType] = GrpcMetadata.ContentTypes.ApplicationJson - } - } - }); - - await call.RequestStream.CompleteAsync(); - try - { - var appendResp = await call.ResponseAsync; - switch (appendResp.ResultCase) - { - case AppendResp.ResultOneofCase.Success: - return; - case AppendResp.ResultOneofCase.WrongExpectedVersion: - throw new WrongExpectedVersionException(); - } - } - catch (RpcException ex) when (ex.Status.StatusCode == StatusCode.DeadlineExceeded) - { - throw new CommitTimeoutException(); - } - } - - private async Task> ReadAllEvents(IPEndPoint endpoint) - { - using var channel = GrpcChannel.ForAddress(new Uri($"https://{endpoint}"), - new GrpcChannelOptions { HttpClient = _httpClient }); - var streamClient = new Streams.StreamsClient(channel); - - using var call = streamClient.Read(new ReadReq - { - Options = new() - { - All = new() - { - Start = new Empty() - }, - ReadDirection = ReadReq.Types.Options.Types.ReadDirection.Forwards, - Count = ulong.MaxValue, - NoFilter = new Empty(), - UuidOption = new() { Structured = new() } - } - }, GetCallOptions()); - - return await (from response in call.ResponseStream.ReadAllAsync() where response.Event != null select response.Event.Event).ToListAsync(); - } - - private MiniClusterNode CreateNode(int index, Endpoints endpoints, EndPoint[] gossipSeeds, - int nodePriority, string intHostAdvertiseAs) => new( - PathName, index, endpoints.InternalTcp, - endpoints.ExternalTcp, endpoints.HttpEndPoint, - subsystems: Array.Empty(), gossipSeeds: gossipSeeds, - nodePriority: nodePriority, intHostAdvertiseAs: intHostAdvertiseAs); - - private Task StartNode(int i, int priority, string intHostAdvertiseAs = null) - { - _nodes[i] = CreateNode(i, _nodeEndpoints[i], _nodeGossipSeeds[i], priority, intHostAdvertiseAs); - _nodes[i].Start(); - return Task.CompletedTask; - } - - private async Task IsNodeReady(IPEndPoint httpEndPoint) - { - var response = await _httpClient.GetAsync($"https://{httpEndPoint}/-/readiness"); - return response.IsSuccessStatusCode; - } - - private async Task WaitForAllNodesToBeLive() - { - for (int i = 0; i < 3; i++) - { - await WaitForNodeToBeLive(i); - } - } - - private async Task WaitForNodeToBeLive(int idx) - { - while (!await IsNodeReady(_nodes[idx].HttpEndPoint)) - { - await Task.Delay(100); - } - } - - private async Task WaitForAllNodesToBeCaughtUp(int maxIdx = 3) - { - while (true) - { - var prevWriter = long.MinValue; - var prevChaser = long.MinValue; - var caughtUp = true; - for (int i = 0; i < maxIdx; i++) - { - var writer = _nodes[i].Db.Config.WriterCheckpoint.ReadNonFlushed(); - var chaser = _nodes[i].Db.Config.ChaserCheckpoint.ReadNonFlushed(); - - if (prevWriter == long.MinValue) - { - prevWriter = writer; - } - - if (prevChaser == long.MinValue) - { - prevChaser = chaser; - } - - if (chaser != writer || writer != prevWriter) - { - caughtUp = false; - } - } - - if (caughtUp) - { - break; - } - - await Task.Delay(100); - } - } - - private async Task ShutdownNode(int i, bool keepDb) => await _nodes[i].Shutdown(keepDb: keepDb); - - private async Task ShutdownAllNodes(int maxIdx = 3, bool keepDb = false) - { - for (int i = 0; i < maxIdx; i++) - { - await ShutdownNode(i, keepDb); - } - } - - private async Task ResignLeader(int leaderIdx) - { - var httpEndPoint = _nodes[leaderIdx].HttpEndPoint; - using var channel = GrpcChannel.ForAddress(new Uri($"https://{httpEndPoint}"), - new GrpcChannelOptions { HttpClient = _httpClient }); - await channel.CreateCallInvoker().AsyncUnaryCall( - ResignNodeMethod, - null, - GetCallOptions(), - new Empty()); - - var start = DateTime.UtcNow; - while (_nodes[leaderIdx].NodeState != VNodeState.Unknown && DateTime.UtcNow - start < TimeSpan.FromSeconds(2)) - { - await Task.Delay(100); - } - } - - [SetUp] - public async Task SetUp() - { - // reset the node states between tests - for (int i = 0; i < 3; i++) - { - await _nodes[i].Shutdown(keepDb: false); - } - - for (int i = 0; i < 3; i++) - { - _nodes[i] = CreateNode(i, _nodeEndpoints[i], _nodeGossipSeeds[i], 0, null); - _nodes[i].Start(); - } - } - - [TestCase(true)] - [TestCase(false)] - [Explicit, Category("LongRunning"), Timeout(80000), NonParallelizable] - public async Task new_events_should_have_correct_event_numbers(bool appendInitialEvent) - { - await WaitForAllNodesToBeLive(); - - if (appendInitialEvent) - { - // append event 0@test - await AppendEvent(_nodes[0].HttpEndPoint, "test", ExpectedVersion.NoStream); - } - - await WaitForAllNodesToBeCaughtUp(); - await ShutdownAllNodes(keepDb: true); - - // make node 1 become the leader by setting its priority to 1 - // node 0 can't become a follower since it can't replicate over internal TCP due to the fake --int-host-advertise-as - await StartNode(0, priority: 0, intHostAdvertiseAs: FakeHostAdvertiseAs); - await StartNode(1, priority: 1, intHostAdvertiseAs: FakeHostAdvertiseAs); - - try - { - await WaitForNodeToBeLive(1).WithTimeout(TimeSpan.FromSeconds(10)); - Assert.AreEqual(VNodeState.Leader, _nodes[1].NodeState); - } - catch (TimeoutException) - { - // want to get stuck in preleader since replication isn't possible - Assert.AreEqual(VNodeState.PreLeader, _nodes[1].NodeState); - return; - } - - // append event 1@test. Expect a commit timeout since there is no quorum. - Assert.ThrowsAsync(async () => - { - await AppendEvent(_nodes[1].HttpEndPoint, "test", appendInitialEvent ? 0 : ExpectedVersion.NoStream); - }); - - // resign the leader node so that it goes into the Unknown state and to trigger new elections - await ResignLeader(1); - - // wait for the node to become Leader again - while (_nodes[1].NodeState != VNodeState.Leader) - { - await Task.Delay(100); - } - - // append event 1@test again. Expect a commit timeout since there is no quorum. - Assert.ThrowsAsync(async () => - { - await AppendEvent(_nodes[1].HttpEndPoint, "test", appendInitialEvent ? 0 : ExpectedVersion.NoStream); - }); - - // shut down both nodes - await ShutdownAllNodes(maxIdx: 2, keepDb: true); - - // start both nodes again without the fake --int-host-advertise-as so that they can form a cluster - await StartNode(0, priority: 0); - await StartNode(1, priority: 1); - try - { - await _nodes[0].Started.WithTimeout(TimeSpan.FromSeconds(10)); - await _nodes[1].Started.WithTimeout(TimeSpan.FromSeconds(10)); - } - catch (TimeoutException) - { - // this is expected in logv3 because by creating the duplicate events above we also - // created duplicate stream records which it will detect and complain about it - throw new Exception($"Couldn't start one or more nodes: {_nodes[0].NodeState} {_nodes[1].NodeState}"); - } - Assert.AreEqual(VNodeState.Follower, _nodes[0].NodeState); - Assert.AreEqual(VNodeState.Leader, _nodes[1].NodeState); - - // wait for data replication - await WaitForAllNodesToBeCaughtUp(maxIdx: 2); - - // read "test" events from $all - var events = - (await ReadAllEvents(_nodes[1].HttpEndPoint)) - .Where(x => x.StreamIdentifier!.StreamName.ToStringUtf8() == "test").ToArray(); - - Assert.AreEqual(appendInitialEvent ? 3 : 2, events.Length); - for (int i = 0; i < (appendInitialEvent ? 3 : 2); i++) - { - Assert.AreEqual(i, events[i].StreamRevision, $"i = {i}, revision = {events[i].StreamRevision}"); - } - } -} 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()); - } -}