diff --git a/src/EventStore.Core.Tests/Authorization/LegacyPolicyVerification.cs b/src/EventStore.Core.Tests/Authorization/LegacyPolicyVerification.cs index 5eb5403acc..76ae0b674c 100644 --- a/src/EventStore.Core.Tests/Authorization/LegacyPolicyVerification.cs +++ b/src/EventStore.Core.Tests/Authorization/LegacyPolicyVerification.cs @@ -497,7 +497,6 @@ public static IEnumerable PolicyTests() yield return CreateOperation(Operations.Node.Options); yield return CreateOperation(Operations.Node.Statistics.Read); yield return CreateOperation(Operations.Node.Statistics.Replication); - yield return CreateOperation(Operations.Node.Statistics.Tcp); yield return CreateOperation(Operations.Node.Statistics.Custom); } diff --git a/src/EventStore.Core.Tests/ClientOperations/specification_with_bare_vnode.cs b/src/EventStore.Core.Tests/ClientOperations/specification_with_bare_vnode.cs index 12331d412f..806e90125d 100644 --- a/src/EventStore.Core.Tests/ClientOperations/specification_with_bare_vnode.cs +++ b/src/EventStore.Core.Tests/ClientOperations/specification_with_bare_vnode.cs @@ -9,7 +9,7 @@ using EventStore.Core.Bus; using EventStore.Core.Certificates; using EventStore.Core.Messaging; -using EventStore.Core.Tests.Services.Transport.Tcp; +using EventStore.Core.Tests.Helpers; using Microsoft.AspNetCore.Builder; namespace EventStore.Core.Tests.ClientOperations; @@ -26,8 +26,8 @@ public void CreateTestNode() var options = new ClusterVNodeOptions() .ReduceMemoryUsageForTests() .RunOnDisk(_dbPath) - .Secure(new X509Certificate2Collection(ssl_connections.GetRootCertificate()), - ssl_connections.GetServerCertificate()); + .Secure(new X509Certificate2Collection(TestCertificates.GetRootCertificate()), + TestCertificates.GetServerCertificate()); _node = new ClusterVNode(options, logFormatFactory, new AuthenticationProviderFactory( c => new InternalAuthenticationProviderFactory(c, options.DefaultUser)), diff --git a/src/EventStore.Core.Tests/ClientOperations/when_committing_a_transaction_with_data.cs b/src/EventStore.Core.Tests/ClientOperations/when_committing_a_transaction_with_data.cs index 94d3acd5a1..d4e1311df3 100644 --- a/src/EventStore.Core.Tests/ClientOperations/when_committing_a_transaction_with_data.cs +++ b/src/EventStore.Core.Tests/ClientOperations/when_committing_a_transaction_with_data.cs @@ -1,7 +1,7 @@ using System; using System.Collections.Generic; using System.Threading; -using EventStore.ClientAPI.Common.Utils; +using EventStore.Common.Utils; using EventStore.Core.Data; using EventStore.Core.Messages; using EventStore.Core.Messaging; diff --git a/src/EventStore.Core.Tests/Helpers/TestCertificates.cs b/src/EventStore.Core.Tests/Helpers/TestCertificates.cs new file mode 100644 index 0000000000..2014a19fb9 --- /dev/null +++ b/src/EventStore.Core.Tests/Helpers/TestCertificates.cs @@ -0,0 +1,66 @@ +using System; +using System.Net; +using System.Security.Cryptography; +using System.Security.Cryptography.X509Certificates; + +namespace EventStore.Core.Tests.Helpers; + +public static class TestCertificates +{ + private static readonly X509Certificate2 Root = CreateRootCertificate("Test Root CA"); + private static readonly X509Certificate2 Server = CreateServerCertificate("localhost", Root); + private static readonly X509Certificate2 OtherServer = CreateServerCertificate("other-node", Root); + private static readonly X509Certificate2 UntrustedRoot = CreateRootCertificate("Untrusted Test Root CA"); + private static readonly X509Certificate2 Untrusted = CreateServerCertificate("untrusted", UntrustedRoot); + + public static X509Certificate2 GetRootCertificate() => + X509CertificateLoader.LoadCertificate(Root.Export(X509ContentType.Cert)); + + public static X509Certificate2 GetServerCertificate() => CloneWithPrivateKey(Server); + + public static X509Certificate2 GetOtherServerCertificate() => CloneWithPrivateKey(OtherServer); + + public static X509Certificate2 GetUntrustedCertificate() => CloneWithPrivateKey(Untrusted); + + private static X509Certificate2 CreateRootCertificate(string commonName) + { + using var key = RSA.Create(2048); + var request = new CertificateRequest( + $"CN={commonName}", key, HashAlgorithmName.SHA256, RSASignaturePadding.Pkcs1); + request.CertificateExtensions.Add(new X509BasicConstraintsExtension(true, false, 0, true)); + request.CertificateExtensions.Add(new X509KeyUsageExtension( + X509KeyUsageFlags.KeyCertSign | X509KeyUsageFlags.CrlSign, true)); + request.CertificateExtensions.Add(new X509SubjectKeyIdentifierExtension(request.PublicKey, false)); + + using var certificate = request.CreateSelfSigned( + DateTimeOffset.UtcNow.AddDays(-1), DateTimeOffset.UtcNow.AddYears(1)); + return CloneWithPrivateKey(certificate); + } + + private static X509Certificate2 CreateServerCertificate(string commonName, X509Certificate2 issuer) + { + using var key = RSA.Create(2048); + var request = new CertificateRequest( + $"CN={commonName}", key, HashAlgorithmName.SHA256, RSASignaturePadding.Pkcs1); + request.CertificateExtensions.Add(new X509BasicConstraintsExtension(false, false, 0, true)); + request.CertificateExtensions.Add(new X509KeyUsageExtension( + X509KeyUsageFlags.DigitalSignature | X509KeyUsageFlags.KeyEncipherment, true)); + request.CertificateExtensions.Add(new X509EnhancedKeyUsageExtension( + [new Oid("1.3.6.1.5.5.7.3.1"), new Oid("1.3.6.1.5.5.7.3.2")], true)); + var names = new SubjectAlternativeNameBuilder(); + names.AddDnsName(commonName); + names.AddDnsName("localhost"); + names.AddIpAddress(IPAddress.Loopback); + request.CertificateExtensions.Add(names.Build()); + + var serial = RandomNumberGenerator.GetBytes(16); + using var certificate = request.Create( + issuer, DateTimeOffset.UtcNow.AddDays(-1), DateTimeOffset.UtcNow.AddMonths(6), serial); + using var certificateWithKey = certificate.CopyWithPrivateKey(key); + return CloneWithPrivateKey(certificateWithKey); + } + + private static X509Certificate2 CloneWithPrivateKey(X509Certificate2 certificate) => + X509CertificateLoader.LoadPkcs12( + certificate.Export(X509ContentType.Pkcs12), string.Empty, X509KeyStorageFlags.Exportable); +} diff --git a/src/EventStore.Core.Tests/Helpers/TestFixtureWithExistingEvents.cs b/src/EventStore.Core.Tests/Helpers/TestFixtureWithExistingEvents.cs index 2544afaf57..af29a907ef 100644 --- a/src/EventStore.Core.Tests/Helpers/TestFixtureWithExistingEvents.cs +++ b/src/EventStore.Core.Tests/Helpers/TestFixtureWithExistingEvents.cs @@ -534,7 +534,7 @@ private void ProcessWrite(IEnvelope envelope, Guid correlationId, string stre _streams[streamId] = list; } - if (expectedVersion != EventStore.ClientAPI.ExpectedVersion.Any) + if (expectedVersion != EventStore.Core.Data.ExpectedVersion.Any) { if (expectedVersion != list.Count - 1) { diff --git a/src/EventStore.Core.Tests/Http/HealthChecks/when_performing_a_live_check.cs b/src/EventStore.Core.Tests/Http/HealthChecks/when_performing_a_live_check.cs index 8b64f5953c..27087c5e43 100644 --- a/src/EventStore.Core.Tests/Http/HealthChecks/when_performing_a_live_check.cs +++ b/src/EventStore.Core.Tests/Http/HealthChecks/when_performing_a_live_check.cs @@ -1,8 +1,6 @@ using System; using System.Net.Http; using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.Core.Tests.ClientAPI.Helpers; using EventStore.Core.Tests.Helpers; using Grpc.Health.V1; using Grpc.Net.Client; @@ -128,11 +126,5 @@ private async Task StartNodeAndWaitForReadiness() { await _node.Start(); _nodeStarted = true; - await _node.WaitForTcpEndPoint().WithTimeout(ReadinessTimeout); - - using var connection = await TestConnectionLifecycle.ReconnectUntilReady( - () => TestConnection.CreateMiniNodeClient(_node.TcpEndPoint), - conn => conn.ReadAllEventsForwardAsync(Position.Start, 1, false, DefaultData.AdminCredentials), - ReadinessTimeout); } } diff --git a/src/EventStore.Core.Tests/Integration/authenticated_requests_made_from_a_follower.cs b/src/EventStore.Core.Tests/Integration/authenticated_requests_made_from_a_follower.cs index e90fdf01ab..884b7614db 100644 --- a/src/EventStore.Core.Tests/Integration/authenticated_requests_made_from_a_follower.cs +++ b/src/EventStore.Core.Tests/Integration/authenticated_requests_made_from_a_follower.cs @@ -4,8 +4,6 @@ using System.Threading.Tasks; using EventStore.Client; using EventStore.Client.Streams; -using EventStore.ClientAPI; -using EventStore.ClientAPI.SystemData; using EventStore.Core.Services.Transport.Grpc; using Google.Protobuf; using Grpc.Core; @@ -94,38 +92,4 @@ await call.RequestStream.WriteAsync(new AppendReq public void work() => Assert.AreEqual(StatusCode.OK, _status.StatusCode); } - [TestFixture(typeof(LogFormat.V2), typeof(string))] - public class via_tcp_should : authenticated_requests_made_from_a_follower - { - private Exception _caughtException; - - protected override async Task Given() - { - var node = GetFollowers()[0]; - await Task.WhenAll(node.AdminUserCreated, node.Started); - - using var connection = EventStoreConnection.Create(ConnectionSettings.Create() - .DisableServerCertificateValidation() - .PreferFollowerNode(), - node.ExternalTcpEndPoint); - await connection.ConnectAsync(); - - try - { - await connection.AppendToStreamAsync(ProtectedStream, ExpectedVersion.NoStream, - new UserCredentials("admin", "changeit"), - new EventData(Guid.NewGuid(), "-", false, Array.Empty(), Array.Empty())); - } - catch (Exception ex) - { - _caughtException = ex; - } - - await base.Given(); - } - - [Test] - [Retry(5)] - public void work() => Assert.Null(_caughtException); - } } diff --git a/src/EventStore.Core.Tests/Integration/specification_with_a_single_node.cs b/src/EventStore.Core.Tests/Integration/specification_with_a_single_node.cs index 7d264ab04c..f957746413 100644 --- a/src/EventStore.Core.Tests/Integration/specification_with_a_single_node.cs +++ b/src/EventStore.Core.Tests/Integration/specification_with_a_single_node.cs @@ -1,10 +1,7 @@ using System; using System.IO; -using System.Net; using System.Threading; using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.ClientAPI.SystemData; using EventStore.Core.Bus; using EventStore.Core.Tests.Helpers; using NUnit.Framework; diff --git a/src/EventStore.Core.Tests/Integration/when_a_single_node_is_restarted_multiple_times.cs b/src/EventStore.Core.Tests/Integration/when_a_single_node_is_restarted_multiple_times.cs index b132166315..b7db44cba5 100644 --- a/src/EventStore.Core.Tests/Integration/when_a_single_node_is_restarted_multiple_times.cs +++ b/src/EventStore.Core.Tests/Integration/when_a_single_node_is_restarted_multiple_times.cs @@ -53,7 +53,6 @@ private async Task GetLastEpochId(Guid? previousEpochId) IndexDirectory = GetFilePathFor("epoch-index"), }); - await _node.WaitForTcpEndPoint().WaitAsync(RestartTimeout); var wait = Stopwatch.StartNew(); while (wait.Elapsed < RestartTimeout) { diff --git a/src/EventStore.Core.Tests/Services/PersistentSubscription/PersistentSubscriptionTests.cs b/src/EventStore.Core.Tests/Services/PersistentSubscription/PersistentSubscriptionTests.cs index 9975fd3b17..cd37d16a0a 100644 --- a/src/EventStore.Core.Tests/Services/PersistentSubscription/PersistentSubscriptionTests.cs +++ b/src/EventStore.Core.Tests/Services/PersistentSubscription/PersistentSubscriptionTests.cs @@ -5,8 +5,6 @@ using System.Text; using System.Threading; using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.ClientAPI.Common; using EventStore.Core.Bus; using EventStore.Core.Data; using EventStore.Core.Helpers; @@ -15,10 +13,10 @@ using EventStore.Core.Messages; using EventStore.Core.Messaging; using EventStore.Core.Metrics; +using EventStore.Core.Services; using EventStore.Core.Services.PersistentSubscription; using EventStore.Core.Services.PersistentSubscription.ConsumerStrategy; using EventStore.Core.Services.Storage.ReaderIndex; -using EventStore.Core.Tests.ClientAPI; using EventStore.Core.Tests.Services.Replication; using EventStore.Core.Tests.TransactionLog; using EventStore.Core.TransactionLog.LogRecords; @@ -2606,58 +2604,6 @@ public void retrying_parked_messages_with_stop_at_replays_parkedEvents_until_tha } } -[Ignore("very long test")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class DeadlockTest : SpecificationWithMiniNode -{ - protected override Task Given() - { - _conn = BuildConnection(_node); - return _conn.ConnectAsync(); - } - - protected override Task When() => Task.CompletedTask; - - [Test] - public async Task read_whilst_ack_doesnt_deadlock_with_request_response_dispatcher() - { - var persistentSubscriptionSettings = PersistentSubscriptionSettings.Create().Build(); - var userCredentials = DefaultData.AdminCredentials; - await _conn.CreatePersistentSubscriptionAsync("TestStream", "TestGroup", persistentSubscriptionSettings, - userCredentials); - - const int count = 5000; - await _conn.AppendToStreamAsync("TestStream", ExpectedVersion.Any, CreateEvent().Take(count)); - - - var received = 0; - var manualResetEventSlim = new ManualResetEventSlim(); - var sub1 = _conn.ConnectToPersistentSubscription("TestStream", "TestGroup", (sub, ev) => - { - received++; - if (received == count) - { - manualResetEventSlim.Set(); - } - - return Task.CompletedTask; - }, - (sub, reason, ex) => { }); - Assert.IsTrue(manualResetEventSlim.Wait(TimeSpan.FromSeconds(30)), - "Failed to receive all events in 2 minutes. Assume event store is deadlocked."); - sub1.Stop(TimeSpan.FromSeconds(10)); - _conn.Close(); - } - - private static IEnumerable CreateEvent() - { - while (true) - { - yield return new EventData(Guid.NewGuid(), "testtype", false, new byte[0], new byte[0]); - } - } -} - public class CheckpointingWithSkippedEvents { [Test] diff --git a/src/EventStore.Core.Tests/Services/Replication/LeaderReplication/when_non_tcp_replica_subscribes.cs b/src/EventStore.Core.Tests/Services/Replication/LeaderReplication/when_replica_subscribes.cs similarity index 96% rename from src/EventStore.Core.Tests/Services/Replication/LeaderReplication/when_non_tcp_replica_subscribes.cs rename to src/EventStore.Core.Tests/Services/Replication/LeaderReplication/when_replica_subscribes.cs index 40f010e3e7..b212493732 100644 --- a/src/EventStore.Core.Tests/Services/Replication/LeaderReplication/when_non_tcp_replica_subscribes.cs +++ b/src/EventStore.Core.Tests/Services/Replication/LeaderReplication/when_replica_subscribes.cs @@ -10,7 +10,7 @@ namespace EventStore.Core.Tests.Services.Replication.LeaderReplication; -public class when_non_tcp_replica_subscribes : WithReplicationService +public class when_replica_subscribes : WithReplicationService { private static readonly ReplicationSessionStatistics Statistics = new( SendQueueSize: 7, diff --git a/src/EventStore.Core.Tests/Services/RequestManagement/Service/RequestManagerServiceSpecification.cs b/src/EventStore.Core.Tests/Services/RequestManagement/Service/RequestManagerServiceSpecification.cs index 4a26eb87b9..4c04d8d10e 100644 --- a/src/EventStore.Core.Tests/Services/RequestManagement/Service/RequestManagerServiceSpecification.cs +++ b/src/EventStore.Core.Tests/Services/RequestManagement/Service/RequestManagerServiceSpecification.cs @@ -1,6 +1,6 @@ using System; using System.Collections.Generic; -using EventStore.ClientAPI.Common.Utils; +using EventStore.Common.Utils; using EventStore.Core.Bus; using EventStore.Core.Data; using EventStore.Core.Messages; diff --git a/src/EventStore.Core.Tests/Services/Storage/AllReader/when_reading_all_with_disallowed_streams.cs b/src/EventStore.Core.Tests/Services/Storage/AllReader/when_reading_all_with_disallowed_streams.cs index 4e4b337578..fad1650335 100644 --- a/src/EventStore.Core.Tests/Services/Storage/AllReader/when_reading_all_with_disallowed_streams.cs +++ b/src/EventStore.Core.Tests/Services/Storage/AllReader/when_reading_all_with_disallowed_streams.cs @@ -2,7 +2,6 @@ using System.Linq; using System.Threading; using System.Threading.Tasks; -using EventStore.Client.Messages; using EventStore.Core.Data; using EventStore.Core.Services.Storage.ReaderIndex; using NUnit.Framework; @@ -49,10 +48,7 @@ public async Task should_filter_out_disallowed_streams_when_reading_events_forwa [Test] public async Task should_filter_out_disallowed_streams_when_reading_events_forward_with_event_type_prefix() { - var filter = new Filter( - Filter.Types.FilterContext.EventType, - Filter.Types.FilterType.Prefix, new[] { "event-type" }); - var eventFilter = EventFilter.Get(true, filter); + var eventFilter = EventFilter.EventType.Prefixes(true, new[] { "event-type" }); var result = await ReadIndex.ReadAllEventsForwardFiltered(_forwardReadPos, 10, 10, eventFilter, CancellationToken.None); @@ -65,10 +61,7 @@ public async Task should_filter_out_disallowed_streams_when_reading_events_forwa [Test] public async Task should_filter_out_disallowed_streams_when_reading_events_forward_with_event_type_regex() { - var filter = new Filter( - Filter.Types.FilterContext.EventType, - Filter.Types.FilterType.Regex, new[] { @"^.*event-type-.*$" }); - var eventFilter = EventFilter.Get(true, filter); + var eventFilter = EventFilter.EventType.Regex(true, @"^.*event-type-.*$"); var result = await ReadIndex.ReadAllEventsForwardFiltered(_forwardReadPos, 10, 10, eventFilter, CancellationToken.None); @@ -81,10 +74,7 @@ public async Task should_filter_out_disallowed_streams_when_reading_events_forwa [Test] public async Task should_filter_out_disallowed_streams_when_reading_events_forward_with_stream_id_prefix() { - var filter = new Filter( - Filter.Types.FilterContext.StreamId, - Filter.Types.FilterType.Prefix, new[] { "$persistentsubscripti" }); - var eventFilter = EventFilter.Get(true, filter); + var eventFilter = EventFilter.StreamName.Prefixes(true, new[] { "$persistentsubscripti" }); var result = await ReadIndex.ReadAllEventsForwardFiltered(_forwardReadPos, 10, 10, eventFilter, CancellationToken.None); @@ -96,10 +86,7 @@ public async Task should_filter_out_disallowed_streams_when_reading_events_forwa [Test] public async Task should_filter_out_disallowed_streams_when_reading_events_forward_with_stream_id_regex() { - var filter = new Filter( - Filter.Types.FilterContext.StreamId, - Filter.Types.FilterType.Regex, new[] { @"^.*istentsubsc.*$" }); - var eventFilter = EventFilter.Get(true, filter); + var eventFilter = EventFilter.StreamName.Regex(true, @"^.*istentsubsc.*$"); var result = await ReadIndex.ReadAllEventsForwardFiltered(_forwardReadPos, 10, 10, eventFilter, CancellationToken.None); @@ -122,10 +109,7 @@ public async Task should_filter_out_disallowed_streams_when_reading_events_backw [Test] public async Task should_filter_out_disallowed_streams_when_reading_events_backward_with_event_type_prefix() { - var filter = new Filter( - Filter.Types.FilterContext.EventType, - Filter.Types.FilterType.Prefix, ["event-type"]); - var eventFilter = EventFilter.Get(true, filter); + var eventFilter = EventFilter.EventType.Prefixes(true, ["event-type"]); var result = await ReadIndex.ReadAllEventsBackwardFiltered(_backwardReadPos, 10, 10, eventFilter, @@ -139,10 +123,7 @@ await ReadIndex.ReadAllEventsBackwardFiltered(_backwardReadPos, 10, 10, eventFil [Test] public async Task should_filter_out_disallowed_streams_when_reading_events_backward_with_event_type_regex() { - var filter = new Filter( - Filter.Types.FilterContext.EventType, - Filter.Types.FilterType.Regex, [@"^.*event-type-.*$"]); - var eventFilter = EventFilter.Get(true, filter); + var eventFilter = EventFilter.EventType.Regex(true, @"^.*event-type-.*$"); var result = await ReadIndex.ReadAllEventsBackwardFiltered(_backwardReadPos, 10, 10, eventFilter, @@ -156,10 +137,7 @@ await ReadIndex.ReadAllEventsBackwardFiltered(_backwardReadPos, 10, 10, eventFil [Test] public async Task should_filter_out_disallowed_streams_when_reading_events_backward_with_stream_id_prefix() { - var filter = new Filter( - Filter.Types.FilterContext.StreamId, - Filter.Types.FilterType.Prefix, ["$persistentsubscripti"]); - var eventFilter = EventFilter.Get(true, filter); + var eventFilter = EventFilter.StreamName.Prefixes(true, ["$persistentsubscripti"]); var result = await ReadIndex.ReadAllEventsBackwardFiltered(_backwardReadPos, 10, 10, eventFilter, @@ -172,10 +150,7 @@ await ReadIndex.ReadAllEventsBackwardFiltered(_backwardReadPos, 10, 10, eventFil [Test] public async Task should_filter_out_disallowed_streams_when_reading_events_backward_with_stream_id_regex() { - var filter = new Filter( - Filter.Types.FilterContext.StreamId, - Filter.Types.FilterType.Regex, [@"^.*istentsubsc.*$"]); - var eventFilter = EventFilter.Get(true, filter); + var eventFilter = EventFilter.StreamName.Regex(true, @"^.*istentsubsc.*$"); var result = await ReadIndex.ReadAllEventsBackwardFiltered(_backwardReadPos, 10, 10, eventFilter, diff --git a/src/EventStore.Core.Tests/Services/Storage/AllReader/when_reading_all_with_filtering.cs b/src/EventStore.Core.Tests/Services/Storage/AllReader/when_reading_all_with_filtering.cs index 1bec5d22c5..6ce9098f24 100644 --- a/src/EventStore.Core.Tests/Services/Storage/AllReader/when_reading_all_with_filtering.cs +++ b/src/EventStore.Core.Tests/Services/Storage/AllReader/when_reading_all_with_filtering.cs @@ -1,7 +1,6 @@ using System; using System.Threading; using System.Threading.Tasks; -using EventStore.Client.Messages; using EventStore.Core.Data; using EventStore.Core.Services.Storage.ReaderIndex; using NUnit.Framework; @@ -35,10 +34,7 @@ protected override async ValueTask WriteTestScenario(CancellationToken token) [Test] public async Task should_read_only_events_forward_with_event_type_prefix() { - var filter = new Filter( - Filter.Types.FilterContext.EventType, - Filter.Types.FilterType.Prefix, ["event-type"]); - var eventFilter = EventFilter.Get(true, filter); + var eventFilter = EventFilter.EventType.Prefixes(true, ["event-type"]); var result = await ReadIndex.ReadAllEventsForwardFiltered(_forwardReadPos, 10, 10, eventFilter, CancellationToken.None); Assert.AreEqual(2, result.Records.Count); @@ -47,10 +43,7 @@ public async Task should_read_only_events_forward_with_event_type_prefix() [Test] public async Task should_read_only_events_forward_with_event_type_regex() { - var filter = new Filter( - Filter.Types.FilterContext.EventType, - Filter.Types.FilterType.Regex, [@"^.*other-event.*$"]); - var eventFilter = EventFilter.Get(true, filter); + var eventFilter = EventFilter.EventType.Regex(true, @"^.*other-event.*$"); var result = await ReadIndex.ReadAllEventsForwardFiltered(_forwardReadPos, 10, 10, eventFilter, CancellationToken.None); Assert.AreEqual(2, result.Records.Count); @@ -59,10 +52,7 @@ public async Task should_read_only_events_forward_with_event_type_regex() [Test] public async Task should_read_only_events_forward_with_stream_id_prefix() { - var filter = new Filter( - Filter.Types.FilterContext.StreamId, - Filter.Types.FilterType.Prefix, ["ES2"]); - var eventFilter = EventFilter.Get(true, filter); + var eventFilter = EventFilter.StreamName.Prefixes(true, ["ES2"]); var result = await ReadIndex.ReadAllEventsForwardFiltered(_forwardReadPos, 10, 10, eventFilter, CancellationToken.None); Assert.AreEqual(1, result.Records.Count); @@ -71,10 +61,7 @@ public async Task should_read_only_events_forward_with_stream_id_prefix() [Test] public async Task should_read_only_events_forward_with_stream_id_regex() { - var filter = new Filter( - Filter.Types.FilterContext.StreamId, - Filter.Types.FilterType.Regex, [@"^.*ES2.*$"]); - var eventFilter = EventFilter.Get(true, filter); + var eventFilter = EventFilter.StreamName.Regex(true, @"^.*ES2.*$"); var result = await ReadIndex.ReadAllEventsForwardFiltered(_forwardReadPos, 10, 10, eventFilter, CancellationToken.None); Assert.AreEqual(1, result.Records.Count); @@ -83,10 +70,7 @@ public async Task should_read_only_events_forward_with_stream_id_regex() [Test] public async Task should_read_only_events_backward_with_event_type_prefix() { - var filter = new Filter( - Filter.Types.FilterContext.EventType, - Filter.Types.FilterType.Prefix, ["event-type"]); - var eventFilter = EventFilter.Get(true, filter); + var eventFilter = EventFilter.EventType.Prefixes(true, ["event-type"]); var result = await ReadIndex.ReadAllEventsBackwardFiltered(_backwardReadPos, 10, 10, eventFilter, @@ -97,10 +81,7 @@ await ReadIndex.ReadAllEventsBackwardFiltered(_backwardReadPos, 10, 10, eventFil [Test] public async Task should_read_only_events_backward_with_event_type_regex() { - var filter = new Filter( - Filter.Types.FilterContext.EventType, - Filter.Types.FilterType.Regex, new[] { @"^.*other-event.*$" }); - var eventFilter = EventFilter.Get(true, filter); + var eventFilter = EventFilter.EventType.Regex(true, @"^.*other-event.*$"); var result = await ReadIndex.ReadAllEventsBackwardFiltered(_backwardReadPos, 10, 10, eventFilter, @@ -111,10 +92,7 @@ await ReadIndex.ReadAllEventsBackwardFiltered(_backwardReadPos, 10, 10, eventFil [Test] public async Task should_read_only_events_backward_with_stream_id_prefix() { - var filter = new Filter( - Filter.Types.FilterContext.StreamId, - Filter.Types.FilterType.Prefix, ["ES2"]); - var eventFilter = EventFilter.Get(true, filter); + var eventFilter = EventFilter.StreamName.Prefixes(true, ["ES2"]); var result = await ReadIndex.ReadAllEventsBackwardFiltered(_backwardReadPos, 10, 10, eventFilter, @@ -125,10 +103,7 @@ await ReadIndex.ReadAllEventsBackwardFiltered(_backwardReadPos, 10, 10, eventFil [Test] public async Task should_read_only_events_backward_with_stream_id_regex() { - var filter = new Filter( - Filter.Types.FilterContext.StreamId, - Filter.Types.FilterType.Regex, [@"^.*ES2.*$"]); - var eventFilter = EventFilter.Get(true, filter); + var eventFilter = EventFilter.StreamName.Regex(true, @"^.*ES2.*$"); var result = await ReadIndex.ReadAllEventsBackwardFiltered(_backwardReadPos, 10, 10, eventFilter, diff --git a/src/EventStore.Core.Tests/Services/Storage/HashCollisions/with_hash_collisions.cs b/src/EventStore.Core.Tests/Services/Storage/HashCollisions/with_hash_collisions.cs index e0428f32d8..46f66f7cd5 100644 --- a/src/EventStore.Core.Tests/Services/Storage/HashCollisions/with_hash_collisions.cs +++ b/src/EventStore.Core.Tests/Services/Storage/HashCollisions/with_hash_collisions.cs @@ -12,7 +12,7 @@ using EventStore.Core.TransactionLog; using EventStore.Core.TransactionLog.LogRecords; using NUnit.Framework; -using ExpectedVersion = EventStore.ClientAPI.ExpectedVersion; +using ExpectedVersion = EventStore.Core.Data.ExpectedVersion; namespace EventStore.Core.Tests.Services.Storage.HashCollisions; diff --git a/src/EventStore.Core.Tests/Services/Storage/Scavenge/when_running_a_scavenge_from_storage_scavenger.cs b/src/EventStore.Core.Tests/Services/Storage/Scavenge/when_running_a_scavenge_from_storage_scavenger.cs index aacbdad527..d596e67c6e 100644 --- a/src/EventStore.Core.Tests/Services/Storage/Scavenge/when_running_a_scavenge_from_storage_scavenger.cs +++ b/src/EventStore.Core.Tests/Services/Storage/Scavenge/when_running_a_scavenge_from_storage_scavenger.cs @@ -1,17 +1,22 @@ using System; using System.Collections.Generic; using System.Linq; -using System.Threading; +using System.Text; using System.Threading.Tasks; -using EventStore.ClientAPI; +using EventStore.Client.Streams; using EventStore.Core.Messages; using EventStore.Core.Messaging; using EventStore.Core.Services; using EventStore.Core.Services.UserManagement; -using EventStore.Core.Tests.ClientAPI.Helpers; using EventStore.Core.Tests.Helpers; +using Google.Protobuf; +using Grpc.Core; +using Grpc.Net.Client; using NUnit.Framework; +using GrpcMetadata = EventStore.Core.Services.Transport.Grpc.Constants.Metadata; using ILogger = Serilog.ILogger; +using ReadEvent = EventStore.Client.Streams.ReadResp.Types.ReadEvent; +using StreamsClient = EventStore.Client.Streams.Streams.StreamsClient; namespace EventStore.Core.Tests.Services.Storage.Scavenge; @@ -21,7 +26,7 @@ public class when_running_scavenge_from_storage_scavenger private static readonly ILogger Log = Serilog.Log.ForContext>(); private static readonly TimeSpan Timeout = TimeSpan.FromSeconds(60); private MiniNode _node; - private List _result; + private List _result; public override async Task TestFixtureSetUp() { @@ -52,35 +57,54 @@ public async Task TearDown() public async Task When() { - using (var conn = TestConnection.Create(_node.TcpEndPoint, TcpType.Ssl, DefaultData.AdminCredentials)) + using var channel = GrpcChannel.ForAddress(new UriBuilder { Scheme = Uri.UriSchemeHttps }.Uri, + new GrpcChannelOptions { HttpClient = _node.HttpClient, DisposeHttpClient = false }); + var client = new StreamsClient(channel); + _result = new List(); + var deadline = DateTime.UtcNow + Timeout; + while (_result.Count < 2 && DateTime.UtcNow < deadline) { - await conn.ConnectAsync(); - var countdown = new CountdownEvent(2); - _result = new List(); - - conn.SubscribeToStreamFrom(SystemStreams.ScavengesStream, null, CatchUpSubscriptionSettings.Default, - (x, y) => + _result.Clear(); + using var call = client.Read(new ReadReq + { + Options = new() { - _result.Add(y); - countdown.Signal(); - return Task.CompletedTask; - }, - _ => Log.Information("Processing events started."), - (x, y, z) => { Log.Information("Subscription dropped: {0}, {1}.", y, z); } - ); + Stream = new() + { + StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(SystemStreams.ScavengesStream) }, + Start = new() + }, + ReadDirection = ReadReq.Types.Options.Types.ReadDirection.Forwards, + ResolveLinks = true, + Count = 2, + NoFilter = new(), + UuidOption = new() { Structured = new() } + } + }, new CallOptions(credentials: CallCredentials.FromInterceptor((_, metadata) => + { + metadata.Add("authorization", + $"Basic {Convert.ToBase64String(Encoding.ASCII.GetBytes("admin:changeit"))}"); + return Task.CompletedTask; + }))); + while (await call.ResponseStream.MoveNext(default)) + if (call.ResponseStream.Current.Event is { } resolvedEvent) + _result.Add(resolvedEvent); - if (!countdown.Wait(Timeout)) + if (_result.Count < 2) { - Assert.Fail("Timeout expired while waiting for events."); + Log.Information("Waiting for the scavenge event stream."); + await Task.Delay(100); } } + + Assert.That(_result.Count, Is.GreaterThanOrEqualTo(2), "Timeout expired while waiting for events."); } [Test] public void should_create_scavenge_started_event_on_index_stream() { var scavengeStartedEvent = - _result.FirstOrDefault(x => x.Event.EventType == SystemEventTypes.ScavengeStarted); + _result.FirstOrDefault(x => x.Event.Metadata[GrpcMetadata.Type] == SystemEventTypes.ScavengeStarted); Assert.IsNotNull(scavengeStartedEvent); } @@ -88,7 +112,7 @@ public void should_create_scavenge_started_event_on_index_stream() public void should_create_scavenge_completed_event_on_index_stream() { var scavengeCompletedEvent = - _result.FirstOrDefault(x => x.Event.EventType == SystemEventTypes.ScavengeCompleted); + _result.FirstOrDefault(x => x.Event.Metadata[GrpcMetadata.Type] == SystemEventTypes.ScavengeCompleted); Assert.IsNotNull(scavengeCompletedEvent); } @@ -96,9 +120,9 @@ public void should_create_scavenge_completed_event_on_index_stream() public void should_link_started_and_completed_events_to_the_same_stream() { var scavengeStartedEvent = - _result.FirstOrDefault(x => x.Event.EventType == SystemEventTypes.ScavengeStarted); + _result.FirstOrDefault(x => x.Event.Metadata[GrpcMetadata.Type] == SystemEventTypes.ScavengeStarted); var scavengeCompletedEvent = - _result.FirstOrDefault(x => x.Event.EventType == SystemEventTypes.ScavengeCompleted); - Assert.AreEqual(scavengeStartedEvent.Event.EventStreamId, scavengeCompletedEvent.Event.EventStreamId); + _result.FirstOrDefault(x => x.Event.Metadata[GrpcMetadata.Type] == SystemEventTypes.ScavengeCompleted); + Assert.AreEqual(scavengeStartedEvent.Event.StreamIdentifier, scavengeCompletedEvent.Event.StreamIdentifier); } } diff --git a/src/EventStore.Core.Tests/Services/Transport/Enumerators/Enumerator.AllSubscription.CombinationTests.cs b/src/EventStore.Core.Tests/Services/Transport/Enumerators/Enumerator.AllSubscription.CombinationTests.cs index 716711dc46..3919930028 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Enumerators/Enumerator.AllSubscription.CombinationTests.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Enumerators/Enumerator.AllSubscription.CombinationTests.cs @@ -2,15 +2,14 @@ using System.Collections.Generic; using System.Text; using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.ClientAPI.Common; +using EventStore.Client.Streams; using EventStore.Core.Data; using EventStore.Core.Services.Transport.Enumerators; using EventStore.Core.Services.UserManagement; using NUnit.Framework; -using ExpectedVersion = EventStore.ClientAPI.ExpectedVersion; using Position = EventStore.Core.Services.Transport.Common.Position; -using ResolvedEvent = EventStore.ClientAPI.ResolvedEvent; +using RecordedEvent = EventStore.Client.Streams.ReadResp.Types.ReadEvent.Types.RecordedEvent; +using SystemStreams = EventStore.Core.Services.SystemStreams; namespace EventStore.Core.Tests.Services.Transport.Enumerators; @@ -120,7 +119,7 @@ public enum CheckpointType private readonly string _streamPrefix; private int _streamSuffix; - private List _events = new(); + private List _events = new(); private int _nextEventIndex; @@ -130,19 +129,18 @@ public AllSubscriptionCombinationTests(TestData testData) _streamPrefix = $"stream-{Guid.NewGuid()}-"; } - private Position GetPosition(ResolvedEvent @event) + private static Position GetPosition(RecordedEvent @event) { - var pos = @event.OriginalPosition!.Value; - return Position.FromInt64(pos.CommitPosition, pos.PreparePosition); + return new Position(@event.CommitPosition, @event.PreparePosition); } - private Position GetPositionPlusOneByte(ResolvedEvent @event) + private static Position GetPositionPlusOneByte(RecordedEvent @event) { var pos = GetPosition(@event); return new Position(pos.CommitPosition + 1, pos.PreparePosition + 1); } - private Position GetPositionMinusOneByte(ResolvedEvent @event) + private static Position GetPositionMinusOneByte(RecordedEvent @event) { var pos = GetPosition(@event); return new Position(pos.CommitPosition - 1, pos.PreparePosition - 1); @@ -161,23 +159,14 @@ private async Task PopulateExistingEvents() { _events.Clear(); - var result = await NodeConnection.ReadAllEventsForwardAsync( - position: EventStore.ClientAPI.Position.Start, - maxCount: 1000, - resolveLinkTos: false); - - foreach (var @event in result.Events) - { - _events.Add(@event); - } + _events.AddRange(await ReadAllEvents()); } private async Task WriteEvent(string stream, string eventType, string data, string metadata) { data ??= string.Empty; metadata ??= string.Empty; - var eventData = new EventData(Guid.NewGuid(), eventType, true, Encoding.UTF8.GetBytes(data), Encoding.UTF8.GetBytes(metadata)); - await NodeConnection.AppendToStreamAsync(stream, ExpectedVersion.Any, eventData); + await AppendToStream(stream, eventType, data, metadata); } private Task WriteEvent() => WriteEvent(_streamPrefix + _streamSuffix++, "type", "{}", null); private Task RevokeAccessWithStreamAcl() => WriteEvent(SystemStreams.MetastreamOf(Core.Services.SystemStreams.AllStream), "$metadata", @"{ ""$acl"": { ""$r"": [] } }", null); @@ -275,8 +264,8 @@ private async Task ReadExpectedEvents(EnumeratorWrapper sub, int nextEventI switch (response) { case Event evt: - var evtPos = _events[nextEventIndex++].OriginalPosition!.Value; - var evtTfPos = new TFPos(evtPos.CommitPosition, evtPos.PreparePosition); + var expectedEvent = _events[nextEventIndex++]; + var evtTfPos = new TFPos((long)expectedEvent.CommitPosition, (long)expectedEvent.PreparePosition); Assert.AreEqual(evtTfPos, evt.EventPosition!.Value); break; case FellBehind: diff --git a/src/EventStore.Core.Tests/Services/Transport/Enumerators/Enumerator.AllSubscriptionFiltered.CombinationTests.cs b/src/EventStore.Core.Tests/Services/Transport/Enumerators/Enumerator.AllSubscriptionFiltered.CombinationTests.cs index 2ac02ecdd7..fadea91183 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Enumerators/Enumerator.AllSubscriptionFiltered.CombinationTests.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Enumerators/Enumerator.AllSubscriptionFiltered.CombinationTests.cs @@ -1,18 +1,20 @@ using System; using System.Collections.Generic; +using System.Linq; using System.Text; +using System.Text.RegularExpressions; using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.ClientAPI.Common; -using EventStore.ClientAPI.Messages; +using EventStore.Client.Streams; using EventStore.Core.Data; using EventStore.Core.Services.Storage.ReaderIndex; using EventStore.Core.Services.Transport.Enumerators; +using EventStore.Core.Services.Transport.Grpc; using EventStore.Core.Services.UserManagement; using NUnit.Framework; -using ExpectedVersion = EventStore.ClientAPI.ExpectedVersion; +using GrpcMetadata = EventStore.Core.Services.Transport.Grpc.Constants.Metadata; using Position = EventStore.Core.Services.Transport.Common.Position; -using ResolvedEvent = EventStore.ClientAPI.ResolvedEvent; +using RecordedEvent = EventStore.Client.Streams.ReadResp.Types.ReadEvent.Types.RecordedEvent; +using SystemStreams = EventStore.Core.Services.SystemStreams; namespace EventStore.Core.Tests.Services.Transport.Enumerators; @@ -379,7 +381,7 @@ public enum EventFilterType private readonly Guid _testGuid; private readonly string _streamPrefix; private int _streamSuffix; - private List _events = new(); + private List _events = new(); private int _nextEventIndex; @@ -390,19 +392,18 @@ public AllSubscriptionFilteredCombinationTests(TestData testData) _streamPrefix = $"stream-{_testGuid}-"; } - private Position GetPosition(ResolvedEvent @event) + private static Position GetPosition(RecordedEvent @event) { - var pos = @event.OriginalPosition!.Value; - return Position.FromInt64(pos.CommitPosition, pos.PreparePosition); + return new Position(@event.CommitPosition, @event.PreparePosition); } - private Position GetPositionPlusOneByte(ResolvedEvent @event) + private static Position GetPositionPlusOneByte(RecordedEvent @event) { var pos = GetPosition(@event); return new Position(pos.CommitPosition + 1, pos.PreparePosition + 1); } - private Position GetPositionMinusOneByte(ResolvedEvent @event) + private static Position GetPositionMinusOneByte(RecordedEvent @event) { var pos = GetPosition(@event); return new Position(pos.CommitPosition - 1, pos.PreparePosition - 1); @@ -421,34 +422,27 @@ private async Task PopulateExistingEvents() { _events.Clear(); - var filter = SubscriptionProps.EventFilterType switch + Func filter = SubscriptionProps.EventFilterType switch { - EventFilterType.None => null, - EventFilterType.StreamPrefix => new Filter(ClientMessage.Filter.FilterContext.StreamId, ClientMessage.Filter.FilterType.Prefix, new[] { $"stream-{_testGuid}" }), - EventFilterType.StreamRegex => new Filter(ClientMessage.Filter.FilterContext.StreamId, ClientMessage.Filter.FilterType.Regex, new[] { $"(.*?){_testGuid}(.*?)" }), - EventFilterType.EventTypePrefix => new Filter(ClientMessage.Filter.FilterContext.EventType, ClientMessage.Filter.FilterType.Prefix, new[] { $"type-{_testGuid}" }), - EventFilterType.EventTypeRegex => new Filter(ClientMessage.Filter.FilterContext.EventType, ClientMessage.Filter.FilterType.Regex, new[] { $"(.*?){_testGuid}(.*?)" }), + EventFilterType.None => _ => true, + EventFilterType.StreamPrefix => @event => StreamName(@event).StartsWith($"stream-{_testGuid}", StringComparison.Ordinal), + EventFilterType.StreamRegex => @event => Regex.IsMatch(StreamName(@event), $"(.*?){_testGuid}(.*?)"), + EventFilterType.EventTypePrefix => @event => EventType(@event).StartsWith($"type-{_testGuid}", StringComparison.Ordinal), + EventFilterType.EventTypeRegex => @event => Regex.IsMatch(EventType(@event), $"(.*?){_testGuid}(.*?)"), _ => throw new ArgumentOutOfRangeException() }; - var result = await NodeConnection.FilteredReadAllEventsForwardAsync( - position: EventStore.ClientAPI.Position.Start, - maxCount: 1000, - resolveLinkTos: false, - filter: filter); + _events.AddRange((await ReadAllEvents()).Where(filter)); - foreach (var @event in result.Events) - { - _events.Add(@event); - } + static string StreamName(RecordedEvent @event) => @event.StreamIdentifier.StreamName.ToStringUtf8(); + static string EventType(RecordedEvent @event) => @event.Metadata[GrpcMetadata.Type]; } private async Task WriteEvent(string stream, string eventType, string data, string metadata) { data ??= string.Empty; metadata ??= string.Empty; - var eventData = new EventData(Guid.NewGuid(), eventType, true, Encoding.UTF8.GetBytes(data), Encoding.UTF8.GetBytes(metadata)); - await NodeConnection.AppendToStreamAsync(stream, ExpectedVersion.Any, eventData); + await AppendToStream(stream, eventType, data, metadata); } private Task WriteEvent() => WriteEvent(_streamPrefix + _streamSuffix++ + "-filtered", $"type-{_testGuid}-filtered", "{}", null); private Task RevokeAccessWithStreamAcl() => WriteEvent(SystemStreams.MetastreamOf(Core.Services.SystemStreams.AllStream), "$metadata", @"{ ""$acl"": { ""$r"": [] } }", null); @@ -564,8 +558,8 @@ private async Task ReadExpectedEvents(EnumeratorWrapper sub, int nextEventI switch (response) { case Event evt: - var evtPos = _events[nextEventIndex++].OriginalPosition!.Value; - var evtTfPos = new TFPos(evtPos.CommitPosition, evtPos.PreparePosition); + var expectedEvent = _events[nextEventIndex++]; + var evtTfPos = new TFPos((long)expectedEvent.CommitPosition, (long)expectedEvent.PreparePosition); lastEventOrCheckpointPos = evtTfPos; numEventsSinceLastCheckpoint++; Assert.AreEqual(evtTfPos, evt.EventPosition!.Value); diff --git a/src/EventStore.Core.Tests/Services/Transport/Enumerators/Enumerator.CombinationTests.cs b/src/EventStore.Core.Tests/Services/Transport/Enumerators/Enumerator.CombinationTests.cs index c65db64107..8f0e43bd56 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Enumerators/Enumerator.CombinationTests.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Enumerators/Enumerator.CombinationTests.cs @@ -1,9 +1,19 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.ClientAPI.SystemData; -using EventStore.Core.Tests.ClientAPI.Helpers; +using EventStore.Client; +using EventStore.Client.Streams; +using EventStore.Core.Services.Transport.Grpc; using EventStore.Core.Tests.Helpers; +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.Services.Transport.Enumerators; @@ -12,8 +22,16 @@ public partial class EnumeratorTests [TestFixture] public class TestFixtureWithMiniNodeConnection : SpecificationWithDirectoryPerTestFixture { + private static readonly CallCredentials AdminCredentials = CallCredentials.FromInterceptor((_, metadata) => + { + metadata.Add("authorization", + $"Basic {Convert.ToBase64String(Encoding.ASCII.GetBytes("admin:changeit"))}"); + return Task.CompletedTask; + }); + protected MiniNode Node { get; private set; } - protected IEventStoreConnection NodeConnection { get; private set; } + private GrpcChannel Channel { get; set; } + private Streams.StreamsClient StreamsClient { get; set; } [OneTimeSetUp] public override async Task TestFixtureSetUp() @@ -21,14 +39,135 @@ public override async Task TestFixtureSetUp() await base.TestFixtureSetUp(); Node = new MiniNode(PathName); await Node.Start(); - NodeConnection = TestConnection.To(Node, TcpType.Ssl, new UserCredentials("admin", "changeit")); - await NodeConnection.ConnectAsync(); + await Node.AdminUserCreated; + Channel = GrpcChannel.ForAddress(new Uri($"https://{Node.HttpEndPoint}"), + new GrpcChannelOptions + { + HttpClient = Node.HttpClient, + DisposeHttpClient = false, + }); + StreamsClient = new Streams.StreamsClient(Channel); } + protected async Task AppendToStream( + string stream, + IEnumerable<(string EventType, byte[] Data, byte[] Metadata)> events) + { + using var call = StreamsClient.Append(GetCallOptions()); + await call.RequestStream.WriteAsync(new AppendReq + { + Options = new() + { + Any = new Empty(), + StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(stream) } + } + }); + + foreach (var @event in events) + { + await call.RequestStream.WriteAsync(new AppendReq + { + ProposedMessage = new() + { + Id = Uuid.NewUuid().ToDto(), + Data = ByteString.CopyFrom(@event.Data), + CustomMetadata = ByteString.CopyFrom(@event.Metadata), + Metadata = + { + [GrpcMetadata.Type] = @event.EventType, + [GrpcMetadata.ContentType] = GrpcMetadata.ContentTypes.ApplicationJson + } + } + }); + } + + await call.RequestStream.CompleteAsync(); + var response = await call.ResponseAsync; + Assert.That(response.ResultCase, Is.EqualTo(AppendResp.ResultOneofCase.Success)); + } + + protected Task AppendToStream(string stream, string eventType, string data, string metadata) => + AppendToStream(stream, + [(eventType, Encoding.UTF8.GetBytes(data ?? string.Empty), Encoding.UTF8.GetBytes(metadata ?? string.Empty))]); + + protected async Task> ReadAllEvents() + { + using var call = StreamsClient.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 Empty() } + } + }, GetCallOptions()); + + return await call.ResponseStream.ReadAllAsync() + .Where(response => response.Event is not null) + .Select(response => response.Event.Event) + .ToArrayAsync(); + } + + protected async Task DeleteStream(string stream, bool hardDelete) + { + if (hardDelete) + { + using var call = StreamsClient.TombstoneAsync(new TombstoneReq + { + Options = new() + { + Any = new Empty(), + StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(stream) } + } + }, GetCallOptions()); + await call.ResponseAsync; + return; + } + + using var deleteCall = StreamsClient.DeleteAsync(new DeleteReq + { + Options = new() + { + Any = new Empty(), + StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(stream) } + } + }, GetCallOptions()); + await deleteCall.ResponseAsync; + } + + protected async Task ReadLastStreamRevision(string stream) + { + using var call = StreamsClient.Read(new ReadReq + { + Options = new() + { + Stream = new() + { + StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(stream) }, + End = new Empty() + }, + ReadDirection = ReadReq.Types.Options.Types.ReadDirection.Backwards, + Count = 1, + NoFilter = new Empty(), + UuidOption = new() { Structured = new Empty() } + } + }, GetCallOptions()); + + var response = await call.ResponseStream.ReadAllAsync() + .FirstAsync(response => response.Event is not null); + return checked((long)response.Event.Event.StreamRevision); + } + + private static CallOptions GetCallOptions() => new( + credentials: AdminCredentials, + deadline: DateTime.UtcNow.AddSeconds(30)); + [OneTimeTearDown] public override async Task TestFixtureTearDown() { - NodeConnection?.Dispose(); + Channel?.Dispose(); await Node.Shutdown(); await base.TestFixtureTearDown(); } diff --git a/src/EventStore.Core.Tests/Services/Transport/Enumerators/Enumerator.StreamSubscription.CombinationTests.cs b/src/EventStore.Core.Tests/Services/Transport/Enumerators/Enumerator.StreamSubscription.CombinationTests.cs index c070aac182..38fe6ab9ac 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Enumerators/Enumerator.StreamSubscription.CombinationTests.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Enumerators/Enumerator.StreamSubscription.CombinationTests.cs @@ -1,14 +1,13 @@ using System; +using System.Linq; using System.Text; using System.Threading.Tasks; -using EventStore.ClientAPI; using EventStore.Core.Data; using EventStore.Core.Services; using EventStore.Core.Services.Transport.Common; using EventStore.Core.Services.Transport.Enumerators; using EventStore.Core.Services.UserManagement; using NUnit.Framework; -using ExpectedVersion = EventStore.ClientAPI.ExpectedVersion; namespace EventStore.Core.Tests.Services.Transport.Enumerators; @@ -678,25 +677,20 @@ private async Task WriteExistingEvents() private async Task WriteEvents(int count) { - var events = new EventData[count]; - for (var i = 0; i < count; i++) - { - events[i] = new EventData(Guid.NewGuid(), "type", true, "{}"u8.ToArray(), Array.Empty()); - } - - await NodeConnection.AppendToStreamAsync(_stream, ExpectedVersion.Any, events); + await AppendToStream(_stream, + Enumerable.Range(0, count) + .Select(_ => ("type", "{}"u8.ToArray(), Array.Empty()))); } private async Task WriteEvent(string stream, string eventType, string data, string metadata) { data ??= string.Empty; metadata ??= string.Empty; - var eventData = new EventData(Guid.NewGuid(), eventType, true, Encoding.UTF8.GetBytes(data), Encoding.UTF8.GetBytes(metadata)); - await NodeConnection.AppendToStreamAsync(stream, ExpectedVersion.Any, eventData); + await AppendToStream(stream, eventType, data, metadata); } private Task WriteEvent() => WriteEvent(_stream, "type", "{}", null); - private async Task SoftDelete() => await NodeConnection.DeleteStreamAsync(_stream, Data.ExpectedVersion.Any, hardDelete: false); - private async Task Tombstone() => await NodeConnection.DeleteStreamAsync(_stream, Data.ExpectedVersion.Any, hardDelete: true); + private Task SoftDelete() => DeleteStream(_stream, hardDelete: false); + private Task Tombstone() => DeleteStream(_stream, hardDelete: true); private async Task RevokeAccessWithStreamAcl() => await WriteEvent(SystemStreams.MetastreamOf(_stream), "$metadata", @"{ ""$acl"": { ""$r"": [] } }", null); private async Task RevokeAccessWithDefaultAcl() => await WriteEvent(SystemStreams.SettingsStream, "update-default-acl", @"{ ""$userStreamAcl"" : { ""$r"" : [] } }", null); @@ -825,9 +819,7 @@ private async Task SetUpForEphemeralStream() return; } - var readResult = await NodeConnection.ReadStreamEventsBackwardAsync(_stream, -1, 1, resolveLinkTos: false); - - _ephemeralStreamLastEventNumber = readResult.LastEventNumber; + _ephemeralStreamLastEventNumber = await ReadLastStreamRevision(_stream); _nextEventNumber = CalculateNextEventNumberFromCheckpoint(); } diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/DeleteTests.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/DeleteTests.cs index 78bc3f2059..efd40bfb12 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/DeleteTests.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/DeleteTests.cs @@ -2,7 +2,6 @@ using System.Linq; using System.Threading.Tasks; using EventStore.Client.Streams; -using EventStore.ClientAPI; using EventStore.Core.Services; using EventStore.Core.Services.Transport.Common; using EventStore.Core.Services.Transport.Grpc; @@ -207,7 +206,7 @@ await call.RequestStream.WriteAsync(new() { Metadata.Type, SystemEventTypes.StreamMetadata }, { Metadata.ContentType, Metadata.ContentTypes.ApplicationJson } }, - Data = ByteString.CopyFromUtf8(StreamMetadata.Build().Build().AsJsonString()) + Data = ByteString.CopyFromUtf8("{}") } }); diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/ReadStreamsForwardTests.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/ReadStreamsForwardTests.cs index eee159ad67..0c7335aab5 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/ReadStreamsForwardTests.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/ReadStreamsForwardTests.cs @@ -1,9 +1,9 @@ using System; using System.Collections.Generic; using System.Linq; +using System.Text; using System.Threading.Tasks; using EventStore.Client.Streams; -using EventStore.ClientAPI; using EventStore.Core.Services.Transport.Grpc; using Google.Protobuf; using Grpc.Core; @@ -60,8 +60,7 @@ await call.RequestStream.WriteAsync(new() [GrpcMetadata.Type] = "$metadata", [GrpcMetadata.ContentType] = GrpcMetadata.ContentTypes.ApplicationJson }, - Data = ByteString.CopyFrom(StreamMetadata.Build().SetTruncateBefore(81).Build() - .AsJsonBytes()) + Data = ByteString.CopyFrom(Encoding.UTF8.GetBytes("{\"$tb\":81}")) } }); diff --git a/src/EventStore.Core.Tests/Services/Transport/Http/Authorization/authorization_tests.cs b/src/EventStore.Core.Tests/Services/Transport/Http/Authorization/authorization_tests.cs index 97d8bfd2f2..788c296bd3 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Http/Authorization/authorization_tests.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Http/Authorization/authorization_tests.cs @@ -5,10 +5,8 @@ using System.Net.Http.Headers; using System.Threading.Tasks; using EventStore.Client.Users; -using EventStore.ClientAPI; using EventStore.Common.Utils; using EventStore.Core.Services; -using EventStore.Core.Tests.ClientAPI.Helpers; using EventStore.Core.Tests.Helpers; using Grpc.Core; using Grpc.Net.Client; @@ -146,12 +144,6 @@ public override async Task TestFixtureSetUp() await base.TestFixtureSetUp(); _node = new MiniNode(PathName); await _node.Start(); - await _node.WaitForTcpEndPoint().WithTimeout(ReadinessTimeout); - - using var connection = await TestConnectionLifecycle.ReconnectUntilReady( - () => TestConnection.CreateMiniNodeClient(_node.TcpEndPoint), - conn => conn.ReadAllEventsForwardAsync(Position.Start, 1, false, DefaultData.AdminCredentials), - ReadinessTimeout); _httpClients["Admin"] = CreateHttpClient("admin", "changeit"); _httpClients["Ops"] = CreateHttpClient("ops", "changeit"); diff --git a/src/EventStore.Core.Tests/Services/UserManagementService/user_management_service.cs b/src/EventStore.Core.Tests/Services/UserManagementService/user_management_service.cs index f933cb9cfc..f6c400ee33 100644 --- a/src/EventStore.Core.Tests/Services/UserManagementService/user_management_service.cs +++ b/src/EventStore.Core.Tests/Services/UserManagementService/user_management_service.cs @@ -2,7 +2,7 @@ using System.Linq; using System.Security.Claims; using System.Threading.Tasks; -using EventStore.ClientAPI.Common.Utils; +using EventStore.Common.Utils; using EventStore.Core.Messages; using EventStore.Core.Services; using EventStore.Core.Services.UserManagement; diff --git a/src/EventStore.Core.Tests/Services/VNode/startup_should.cs b/src/EventStore.Core.Tests/Services/VNode/startup_should.cs index 32c20dac8e..8092cda91f 100644 --- a/src/EventStore.Core.Tests/Services/VNode/startup_should.cs +++ b/src/EventStore.Core.Tests/Services/VNode/startup_should.cs @@ -12,7 +12,6 @@ using EventStore.Core.Configuration.Sources; using EventStore.Core.Services.Monitoring; using EventStore.Core.Tests.Helpers; -using EventStore.Core.Tests.Services.Transport.Tcp; using Microsoft.AspNetCore.Hosting; using Microsoft.AspNetCore.TestHost; using Microsoft.Extensions.Configuration; @@ -31,8 +30,6 @@ public async Task propagate_cancellation_into_startup_tasks() var startupTaskStarted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); var blockingStartupTask = new BlockingStartupTask(startupTaskStarted); var ip = IPAddress.Loopback; - var tcpEndPoint = new IPEndPoint(ip, PortsHelper.GetAvailablePort(ip)); - var internalEndPoint = new IPEndPoint(ip, PortsHelper.GetAvailablePort(ip)); var httpEndPoint = new IPEndPoint(ip, PortsHelper.GetAvailablePort(ip)); var options = new ClusterVNodeOptions @@ -44,11 +41,6 @@ public async Task propagate_cancellation_into_startup_tasks() StatsPeriodSec = 60 * 60, WorkerThreads = 1 }, - Interface = new() - { - ReplicationHeartbeatInterval = 10_000, - ReplicationHeartbeatTimeout = 10_000 - }, Cluster = new() { DiscoverViaDns = false, @@ -68,21 +60,12 @@ public async Task propagate_cancellation_into_startup_tasks() LoadedOptions = ClusterVNodeOptions.GetLoadedOptions(new ConfigurationBuilder() .AddEventStoreDefaultValues() .Build()), - }.Secure(new X509Certificate2Collection(ssl_connections.GetRootCertificate()), - ssl_connections.GetServerCertificate()) - .WithReplicationEndpointOn(internalEndPoint) - .WithExternalTcpOn(tcpEndPoint) + }.Secure(new X509Certificate2Collection(TestCertificates.GetRootCertificate()), + TestCertificates.GetServerCertificate()) .WithNodeEndpointOn(httpEndPoint) .RunOnDisk(System.IO.Path.Combine(PathName, "db")); - var configuration = new ConfigurationBuilder() - .AddInMemoryCollection(new KeyValuePair[] { - new("EventStore:TcpUnitTestPlugin:NodeTcpPort", tcpEndPoint.Port.ToString()), - new("EventStore:TcpUnitTestPlugin:NodeHeartbeatInterval", "10000"), - new("EventStore:TcpUnitTestPlugin:NodeHeartbeatTimeout", "10000"), - new("EventStore:TcpUnitTestPlugin:Insecure", options.Application.Insecure.ToString()), - }) - .Build(); + var configuration = new ConfigurationBuilder().Build(); var node = new ClusterVNode( options, diff --git a/src/EventStore.Core.Tests/TransactionLog/Truncation/when_truncating_database.cs b/src/EventStore.Core.Tests/TransactionLog/Truncation/when_truncating_database.cs index a21d8f3546..8047641bf1 100644 --- a/src/EventStore.Core.Tests/TransactionLog/Truncation/when_truncating_database.cs +++ b/src/EventStore.Core.Tests/TransactionLog/Truncation/when_truncating_database.cs @@ -20,7 +20,6 @@ public async Task everything_should_go_fine() var miniNode = new MiniNode(PathName); await miniNode.Start(); - var tcpPort = miniNode.TcpEndPoint.Port; var httpPort = miniNode.HttpEndPoint.Port; const int cnt = 50; var countdown = new CountdownEvent(cnt); @@ -43,12 +42,12 @@ public async Task everything_should_go_fine() await miniNode.Shutdown(keepDb: true); // --- first restart and truncation - miniNode = new MiniNode(PathName, tcpPort, httpPort); + miniNode = new MiniNode(PathName, httpPort); await miniNode.Start(); await miniNode.Shutdown(keepDb: true); // --- second restart after truncation - miniNode = new MiniNode(PathName, tcpPort, httpPort); + miniNode = new MiniNode(PathName, httpPort); await miniNode.Start(); Assert.AreEqual(-1, miniNode.Db.Config.TruncateCheckpoint.Read()); Assert.That(miniNode.Db.Config.WriterCheckpoint.Read(), Is.GreaterThanOrEqualTo(truncatePosition)); @@ -61,7 +60,7 @@ public async Task everything_should_go_fine() await miniNode.Shutdown(keepDb: true); // -- third restart - miniNode = new MiniNode(PathName, tcpPort, httpPort); + miniNode = new MiniNode(PathName, httpPort); Assert.AreEqual(-1, miniNode.Db.Config.TruncateCheckpoint.Read()); await miniNode.Start(); @@ -99,16 +98,14 @@ public async Task with_truncate_position_in_completed_chunk_everything_should_go await miniNode.Shutdown(keepDb: true); - var tcpPort = miniNode.TcpEndPoint.Port; - // --- first restart and truncation - miniNode = new MiniNode(PathName, tcpPort, httpPort, chunkSize: chunkSize, + miniNode = new MiniNode(PathName, httpPort, chunkSize: chunkSize, cachedChunkSize: cachedSize); await miniNode.Start(); await miniNode.Shutdown(keepDb: true); // --- second restart after truncation - miniNode = new MiniNode(PathName, tcpPort, httpPort, chunkSize: chunkSize, + miniNode = new MiniNode(PathName, httpPort, chunkSize: chunkSize, cachedChunkSize: cachedSize); await miniNode.Start(); Assert.AreEqual(-1, miniNode.Db.Config.TruncateCheckpoint.Read()); diff --git a/src/EventStore.Core.Tests/Transforms/TransformTests.cs b/src/EventStore.Core.Tests/Transforms/TransformTests.cs index b760c89303..0f93e821af 100644 --- a/src/EventStore.Core.Tests/Transforms/TransformTests.cs +++ b/src/EventStore.Core.Tests/Transforms/TransformTests.cs @@ -4,8 +4,8 @@ using System.Linq; using System.Threading; using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.Core.Tests.ClientAPI.Helpers; +using EventStore.Client.Streams; +using EventStore.Core.Services.Transport.Grpc; using EventStore.Core.Tests.Helpers; using EventStore.Core.Tests.Transforms.BitFlip; using EventStore.Core.Tests.Transforms.ByteDup; @@ -14,7 +14,11 @@ using EventStore.Core.TransactionLog.Chunks.TFChunk; using EventStore.Core.Transforms.Identity; using EventStore.Plugins.Transforms; +using Google.Protobuf; +using Grpc.Net.Client; using NUnit.Framework; +using GrpcMetadata = EventStore.Core.Services.Transport.Grpc.Constants.Metadata; +using StreamsClient = EventStore.Client.Streams.Streams.StreamsClient; namespace EventStore.Core.Tests.Transforms; @@ -33,7 +37,7 @@ public class TransformTests : SpecificationWithDirectoryP public async Task transform_works(string transform) { MiniNode node = null; - IEventStoreConnection connection = null; + GrpcChannel connection = null; var dbPath = Path.Combine(PathName, $"node-{Guid.NewGuid()}"); try { @@ -82,7 +86,7 @@ private async ValueTask VerifyChecksums(MiniNode node, Ca } } - private async Task<(MiniNode, IEventStoreConnection)> CreateNode(string dbPath, string transform) + private async Task<(MiniNode, GrpcChannel)> CreateNode(string dbPath, string transform) { IDbTransform dbTransform = transform switch { @@ -103,14 +107,13 @@ private async ValueTask VerifyChecksums(MiniNode node, Ca await node.Start(StartupTimeout); var connection = BuildConnection(node); - await connection.ConnectAsync(); return (node, connection); } private static async Task ShutdownNode( MiniNode node, - IEventStoreConnection connection, + GrpcChannel connection, bool keepDb = false) { if (node is not null) @@ -121,52 +124,82 @@ private static async Task ShutdownNode( connection?.Dispose(); } - private static IEventStoreConnection BuildConnection(MiniNode node) + private static GrpcChannel BuildConnection(MiniNode node) { - return TestConnection.Create(node.TcpEndPoint); + return GrpcChannel.ForAddress(new UriBuilder { Scheme = Uri.UriSchemeHttps }.Uri, + new GrpcChannelOptions { HttpClient = node.HttpClient, DisposeHttpClient = false }); } - private static async Task WriteEvents(IEventStoreConnection connection) + private static async Task WriteEvents(GrpcChannel connection) { var writtenIds = new List(); + var client = new StreamsClient(connection); for (var i = 0; i < NumEvents / BatchSize; i++) { var events = CreateEventBatch(BatchSize); - await connection.AppendToStreamAsync("test", ExpectedVersion.Any, events); - writtenIds.AddRange(events.Select(x => x.EventId)); + using var call = client.Append(); + await call.RequestStream.WriteAsync(new AppendReq + { + Options = new() + { + Any = new(), + StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8("test") } + } + }); + foreach (var @event in events) + await call.RequestStream.WriteAsync(new AppendReq { ProposedMessage = @event }); + await call.RequestStream.CompleteAsync(); + await call.ResponseAsync; + writtenIds.AddRange(events.Select(x => Uuid.FromDto(x.Id).ToGuid())); } return writtenIds.ToArray(); } - private static EventData[] CreateEventBatch(int numEvents) + private static AppendReq.Types.ProposedMessage[] CreateEventBatch(int numEvents) { - var events = new EventData[numEvents]; + var events = new AppendReq.Types.ProposedMessage[numEvents]; for (var i = 0; i < numEvents; i++) { - events[i] = new(eventId: Guid.NewGuid(), - type: "testEvent", - isJson: true, - data: "{ \"foo\":\"bar\" }"u8.ToArray(), - metadata: null); + events[i] = new() + { + Id = Uuid.NewUuid().ToDto(), + Data = ByteString.CopyFromUtf8("{ \"foo\":\"bar\" }"), + CustomMetadata = ByteString.Empty, + Metadata = { + { GrpcMetadata.Type, "testEvent" }, + { GrpcMetadata.ContentType, GrpcMetadata.ContentTypes.ApplicationJson } + } + }; } return events; } - private static async Task VerifyEvents(IEventStoreConnection connection, Guid[] writtenIds) + private static async Task VerifyEvents(GrpcChannel connection, Guid[] writtenIds) { - StreamEventsSlice slice; - var nextEventNumber = 0L; - var readIds = new List(); - do + var client = new StreamsClient(connection); + using var call = client.Read(new ReadReq { - slice = await connection.ReadStreamEventsForwardAsync("test", nextEventNumber, BatchSize, resolveLinkTos: false); - readIds.AddRange(slice.Events.Select(evt => evt.Event.EventId)); - nextEventNumber = slice.NextEventNumber; - } while (!slice.IsEndOfStream); + Options = new() + { + Stream = new() + { + StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8("test") }, + Start = new() + }, + ReadDirection = ReadReq.Types.Options.Types.ReadDirection.Forwards, + Count = ulong.MaxValue, + NoFilter = new(), + UuidOption = new() { Structured = new() } + } + }); + var readIds = new List(); + while (await call.ResponseStream.MoveNext(default)) + if (call.ResponseStream.Current.Event is { } resolvedEvent) + readIds.Add(Uuid.FromDto(resolvedEvent.Event.Id).ToGuid()); Assert.AreEqual(NumEvents, readIds.Count); Assert.True(writtenIds.SequenceEqual(readIds)); diff --git a/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/ClusterVNodeOptionsScenarios.cs b/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/ClusterVNodeOptionsScenarios.cs index 7590813f89..a3493c56f4 100644 --- a/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/ClusterVNodeOptionsScenarios.cs +++ b/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/ClusterVNodeOptionsScenarios.cs @@ -8,7 +8,7 @@ using EventStore.Core.Certificates; using EventStore.Core.LogAbstraction; using EventStore.Core.Tests; -using EventStore.Core.Tests.Services.Transport.Tcp; +using EventStore.Core.Tests.Helpers; using NUnit.Framework; namespace EventStore.Core.XUnit.Tests.Configuration.ClusterNodeOptionsTests; @@ -33,8 +33,8 @@ public override async Task TestFixtureSetUp() _options = WithOptions(options .RunOnDisk(PathName) - .Secure(new X509Certificate2Collection(ssl_connections.GetRootCertificate()), - ssl_connections.GetServerCertificate())); + .Secure(new X509Certificate2Collection(TestCertificates.GetRootCertificate()), + TestCertificates.GetServerCertificate())); _node = new ClusterVNode(_options, _logFormatFactory, new AuthenticationProviderFactory(c => new InternalAuthenticationProviderFactory(c, _options.DefaultUser)), @@ -84,8 +84,8 @@ public override async Task TestFixtureSetUp() .ReduceMemoryUsageForTests() .InCluster(_clusterSize) .RunOnDisk(PathName) - .Secure(new X509Certificate2Collection(ssl_connections.GetRootCertificate()), - ssl_connections.GetServerCertificate())); + .Secure(new X509Certificate2Collection(TestCertificates.GetRootCertificate()), + TestCertificates.GetServerCertificate())); _node = new ClusterVNode(_options, _logFormatFactory, new AuthenticationProviderFactory(_ => new InternalAuthenticationProviderFactory(_, _options.DefaultUser)), diff --git a/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_building/with_cluster_node_and_custom_settings.cs b/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_building/with_cluster_node_and_custom_settings.cs index d80eac39fd..2257c1bdbd 100644 --- a/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_building/with_cluster_node_and_custom_settings.cs +++ b/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_building/with_cluster_node_and_custom_settings.cs @@ -127,8 +127,6 @@ protected override ClusterVNodeOptions WithOptions(ClusterVNodeOptions options) { return options .Insecure() - .WithExternalTcpOn(new IPEndPoint(IPAddress.Loopback, 11130)) - .WithReplicationEndpointOn(new IPEndPoint(IPAddress.Loopback, 11120)) .AdvertiseExternalHostAs(new DnsEndPoint("196.168.1.1", 11131)) .AdvertiseNodeAs(new DnsEndPoint("196.168.1.1", 21130)); } @@ -139,12 +137,6 @@ public void should_set_the_custom_advertise_info_for_external() Assert.AreEqual(new DnsEndPoint("196.168.1.1", 21130), _node.GossipAdvertiseInfo.HttpEndPoint); } - - [Test] - public void should_set_the_loopback_address_as_advertise_info_for_internal() - { - Assert.AreEqual(new DnsEndPoint(IPAddress.Loopback.ToString(), 11120), _node.GossipAdvertiseInfo.InternalTcp); - } } [TestFixture(typeof(LogFormat.V2), typeof(string))] @@ -154,8 +146,7 @@ protected override ClusterVNodeOptions WithOptions(ClusterVNodeOptions options) { return options .Insecure() - .WithReplicationEndpointOn(new IPEndPoint(IPAddress.Any, 11120)) - .WithExternalTcpOn(new IPEndPoint(IPAddress.Any, 11130)) + .WithNodeEndpointOn(new IPEndPoint(IPAddress.Any, 2113)) .AdvertiseExternalHostAs(new DnsEndPoint("10.0.0.1", 11131)); } @@ -165,13 +156,6 @@ public void should_set_the_custom_advertise_info_for_external() Assert.AreEqual(new DnsEndPoint("10.0.0.1", 2113), _node.GossipAdvertiseInfo.HttpEndPoint); } - - [Test] - public void should_set_the_non_loopback_address_as_advertise_info_for_internal() - { - Assert.AreEqual(new DnsEndPoint(IPFinder.GetNonLoopbackAddress().ToString(), 11120), - _node.GossipAdvertiseInfo.InternalTcp); - } } [TestFixture(typeof(LogFormat.V2), typeof(string))] @@ -181,9 +165,7 @@ protected override ClusterVNodeOptions WithOptions(ClusterVNodeOptions options) { return options .Insecure() - .WithNodeEndpointOn(new IPEndPoint(IPAddress.Any, 21130)) - .WithExternalTcpOn(new IPEndPoint(IPAddress.Any, 11130)) - .WithReplicationEndpointOn(new IPEndPoint(IPAddress.Loopback, 11120)); + .WithNodeEndpointOn(new IPEndPoint(IPAddress.Any, 21130)); } [Test] @@ -192,12 +174,6 @@ public void should_use_the_non_default_loopback_ip_as_advertise_info_for_externa Assert.AreEqual(new DnsEndPoint(IPFinder.GetNonLoopbackAddress().ToString(), 21130), _node.GossipAdvertiseInfo.HttpEndPoint); } - - [Test] - public void should_use_loopback_ip_as_advertise_info_for_internal() - { - Assert.AreEqual(new DnsEndPoint(IPAddress.Loopback.ToString(), 11120), _node.GossipAdvertiseInfo.InternalTcp); - } } [TestFixture(typeof(LogFormat.V2), typeof(string))] @@ -209,8 +185,6 @@ protected override ClusterVNodeOptions WithOptions(ClusterVNodeOptions options) return options .Insecure() .WithNodeEndpointOn(new IPEndPoint(IPAddress.Loopback, 21130)) - .WithExternalTcpOn(new IPEndPoint(IPAddress.Loopback, 11130)) - .WithReplicationEndpointOn(new IPEndPoint(IPAddress.Any, 11120)) .AdvertiseExternalHostAs(new DnsEndPoint("10.0.0.1", 11131)) .AdvertiseNodeAs(new DnsEndPoint("10.0.0.1", 21131)); } @@ -221,13 +195,6 @@ public void should_set_the_custom_advertise_info_for_external() Assert.AreEqual(new DnsEndPoint("10.0.0.1", 21131), _node.GossipAdvertiseInfo.HttpEndPoint); } - - [Test] - public void should_use_the_non_default_loopback_ip_as_advertise_info_for_internal() - { - Assert.AreEqual(new DnsEndPoint(IPFinder.GetNonLoopbackAddress().ToString(), 11120), - _node.GossipAdvertiseInfo.InternalTcp); - } } [TestFixture(typeof(LogFormat.V2), typeof(string))] diff --git a/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_building/with_single_node_and_custom_settings.cs b/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_building/with_single_node_and_custom_settings.cs index b0c3643ceb..0cd84704a9 100644 --- a/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_building/with_single_node_and_custom_settings.cs +++ b/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_building/with_single_node_and_custom_settings.cs @@ -27,15 +27,10 @@ public void should_set_the_db_path() public class with_custom_ip_endpoints : SingleNodeScenario { private readonly IPEndPoint _httpEndPoint = new(IPAddress.Parse("127.0.1.15"), 1113); - private readonly IPEndPoint _internalTcp = new(IPAddress.Parse("127.0.1.15"), 1114); - private readonly IPEndPoint _externalTcp = new(IPAddress.Parse("127.0.1.15"), 1115); protected override ClusterVNodeOptions WithOptions(ClusterVNodeOptions options) { - return options - .WithNodeEndpointOn(_httpEndPoint) - .WithExternalTcpOn(_externalTcp) - .WithReplicationEndpointOn(_internalTcp); + return options.WithNodeEndpointOn(_httpEndPoint); } [Test] @@ -43,12 +38,6 @@ public void should_set_http_endpoint() { Assert.AreEqual(_httpEndPoint, _node.NodeInfo.HttpEndPoint); } - - [Test] - public void should_set_internal_tcp_endpoint() - { - Assert.AreEqual(_internalTcp, _node.NodeInfo.InternalSecureTcp); - } } [TestFixture(typeof(LogFormat.V2), typeof(string))] @@ -118,10 +107,7 @@ public void should_set_max_chunk_size_to_the_size_of_the_number_of_cached_chunks [TestFixture(typeof(LogFormat.V2), typeof(string))] public class with_custom_advertise_as : SingleNodeScenario { - private readonly IPEndPoint _intTcpEndpoint = new(IPAddress.Parse(InternalIp), 1111); - private readonly IPEndPoint _extTcpEndpoint = new(IPAddress.Parse(ExternalIp), 1113); private readonly IPEndPoint _httpEndpoint = new(IPAddress.Parse(ExternalIp), 1116); - const string InternalIp = "127.0.1.1"; const string ExternalIp = "127.0.1.2"; @@ -129,25 +115,15 @@ protected override ClusterVNodeOptions WithOptions(ClusterVNodeOptions options) { return options .WithNodeEndpointOn(_httpEndpoint) - .WithExternalTcpOn(_extTcpEndpoint) - .WithReplicationEndpointOn(_intTcpEndpoint) - .AdvertiseInternalHostAs(new DnsEndPoint($"{InternalIp}.com", _intTcpEndpoint.Port + 1000)) - .AdvertiseExternalHostAs(new DnsEndPoint($"{ExternalIp}.com", _extTcpEndpoint.Port + 1000)) + .AdvertiseExternalHostAs(new DnsEndPoint($"{ExternalIp}.com", _httpEndpoint.Port + 1000)) .AdvertiseNodeAs(new DnsEndPoint($"{ExternalIp}.com", _httpEndpoint.Port + 1000)); } [Test] public void should_set_the_advertise_as_info_to_the_specified() { - Assert.AreEqual(null, _node.GossipAdvertiseInfo.InternalTcp); - Assert.AreEqual(null, _node.GossipAdvertiseInfo.ExternalTcp); - Assert.AreEqual(new DnsEndPoint($"{InternalIp}.com", _intTcpEndpoint.Port + 1000), - _node.GossipAdvertiseInfo.InternalSecureTcp); Assert.AreEqual(new DnsEndPoint($"{ExternalIp}.com", _httpEndpoint.Port + 1000), _node.GossipAdvertiseInfo.HttpEndPoint); - Assert.AreEqual($"{InternalIp}.com", _node.GossipAdvertiseInfo.AdvertiseInternalHostAs); - Assert.AreEqual($"{ExternalIp}.com", _node.GossipAdvertiseInfo.AdvertiseExternalHostAs); - Assert.AreEqual(_httpEndpoint.Port + 1000, _node.GossipAdvertiseInfo.AdvertiseHttpPortAs); } } diff --git a/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_building/with_secure_tcp.cs b/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_building/with_tls.cs similarity index 70% rename from src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_building/with_secure_tcp.cs rename to src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_building/with_tls.cs index d8c2c18659..a0f6942090 100644 --- a/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_building/with_secure_tcp.cs +++ b/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_building/with_tls.cs @@ -1,6 +1,5 @@ using System; using System.IO; -using System.Net; using System.Security.Cryptography.X509Certificates; using System.Threading.Tasks; using EventStore.Common.Utils; @@ -11,22 +10,18 @@ using EventStore.Core.Authorization.AuthorizationPolicies; using EventStore.Core.Certificates; using EventStore.Core.Tests; -using EventStore.Core.Tests.Services.Transport.Tcp; +using EventStore.Core.Tests.Helpers; using NUnit.Framework; namespace EventStore.Core.XUnit.Tests.Configuration.ClusterNodeOptionsTests.when_building; [Category("LongRunning")] [TestFixture(typeof(LogFormat.V2), typeof(string))] -public class with_ssl_enabled_and_using_a_security_certificate_from_file : SingleNodeScenario +public class with_tls_enabled_and_using_a_security_certificate_from_file : SingleNodeScenario { - private readonly IPEndPoint _internalSecTcp = new(IPAddress.Parse("127.0.1.15"), 1114); - private readonly IPEndPoint _externalSecTcp = new(IPAddress.Parse("127.0.1.15"), 1115); - protected override ClusterVNodeOptions WithOptions(ClusterVNodeOptions options) { - - return options.WithReplicationEndpointOn(_internalSecTcp).WithExternalTcpOn(_externalSecTcp) with + return options with { CertificateFile = new() { @@ -43,16 +38,10 @@ public void should_set_certificate() Assert.AreNotEqual("n/a", _options.Certificate == null ? "n/a" : _options.Certificate.ToString()); } - [Test] - public void should_set_internal_secure_tcp_endpoint() - { - Assert.AreEqual(_internalSecTcp, _node.NodeInfo.InternalSecureTcp); - } - private string GetCertificatePath() { var filePath = Path.Combine(PathName, $"cert-{Guid.NewGuid()}.p12"); - var cert = ssl_connections.GetUntrustedCertificate(); + var cert = TestCertificates.GetUntrustedCertificate(); using var fileStream = File.Create(filePath); fileStream.Write(cert.ExportToPkcs12()); @@ -62,18 +51,13 @@ private string GetCertificatePath() } [TestFixture(typeof(LogFormat.V2), typeof(string))] -public class with_ssl_enabled_and_using_a_security_certificate : SingleNodeScenario +public class with_tls_enabled_and_using_a_security_certificate : SingleNodeScenario { - private readonly IPEndPoint _internalSecTcp = new(IPAddress.Parse("127.0.1.15"), 1114); - private readonly IPEndPoint _externalSecTcp = new(IPAddress.Parse("127.0.1.15"), 1115); - private readonly X509Certificate2 _certificate = ssl_connections.GetServerCertificate(); + private readonly X509Certificate2 _certificate = TestCertificates.GetServerCertificate(); protected override ClusterVNodeOptions WithOptions(ClusterVNodeOptions options) { - return options - .WithReplicationEndpointOn(_internalSecTcp) - .WithExternalTcpOn(_externalSecTcp) - .Secure(new X509Certificate2Collection(ssl_connections.GetRootCertificate()), _certificate); + return options.Secure(new X509Certificate2Collection(TestCertificates.GetRootCertificate()), _certificate); } [Test] @@ -81,16 +65,10 @@ public void should_set_certificate() { Assert.AreNotEqual("n/a", _options.Certificate == null ? "n/a" : _options.Certificate.ToString()); } - - [Test] - public void should_set_internal_secure_tcp_endpoint() - { - Assert.AreEqual(_internalSecTcp, _node.NodeInfo.InternalSecureTcp); - } } [TestFixture(typeof(LogFormat.V2), typeof(string))] -public class with_secure_tcp_endpoints_and_no_certificates : SpecificationWithDirectoryPerTestFixture +public class with_tls_enabled_and_no_certificates : SpecificationWithDirectoryPerTestFixture { private ClusterVNodeOptions _options; private Exception _caughtException; @@ -98,14 +76,9 @@ public class with_secure_tcp_endpoints_and_no_certificates(_options, LogFormatHelper.LogFormatFactory, @@ -168,7 +141,6 @@ public void should_disable_transport_tls() { Assert.IsTrue(_node.DisableHttps); Assert.IsTrue(_options.Application.TlsDisabled()); - Assert.AreEqual(new IPEndPoint(IPAddress.Loopback, 1112), _node.NodeInfo.InternalTcp); } [Test] diff --git a/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_shutting_down_an_isolated_cluster_member.cs b/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_shutting_down_an_isolated_cluster_member.cs index 900aeddce5..d496a8701f 100644 --- a/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_shutting_down_an_isolated_cluster_member.cs +++ b/src/EventStore.Core.XUnit.Tests/Configuration/ClusterNodeOptionsTests/when_shutting_down_an_isolated_cluster_member.cs @@ -13,7 +13,7 @@ using EventStore.Core.LogAbstraction; using EventStore.Core.Messages; using EventStore.Core.Tests; -using EventStore.Core.Tests.Services.Transport.Tcp; +using EventStore.Core.Tests.Helpers; using NUnit.Framework; namespace EventStore.Core.XUnit.Tests.Configuration.ClusterNodeOptionsTests; @@ -43,8 +43,8 @@ static string[] Snapshot(HashSet services) .ReduceMemoryUsageForTests() .InCluster(3) .RunOnDisk(PathName) - .Secure(new X509Certificate2Collection(ssl_connections.GetRootCertificate()), - ssl_connections.GetServerCertificate()); + .Secure(new X509Certificate2Collection(TestCertificates.GetRootCertificate()), + TestCertificates.GetServerCertificate()); var node = new ClusterVNode(options, logFormatFactory, new AuthenticationProviderFactory(c =>