diff --git a/src/EventStore.Projections.Core.Tests/ClientAPI/Cluster/specification_with_standard_projections_runnning.cs b/src/EventStore.Projections.Core.Tests/ClientAPI/Cluster/specification_with_standard_projections_runnning.cs index 68e0e66259..7e379c45be 100644 --- a/src/EventStore.Projections.Core.Tests/ClientAPI/Cluster/specification_with_standard_projections_runnning.cs +++ b/src/EventStore.Projections.Core.Tests/ClientAPI/Cluster/specification_with_standard_projections_runnning.cs @@ -1,32 +1,45 @@ using System; using System.Collections.Generic; -using System.Diagnostics; using System.Linq; using System.Net; +using System.Net.Http; using System.Text; +using System.Threading; using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.ClientAPI.SystemData; +using EventStore.Client; +using EventStore.Client.Streams; using EventStore.Common.Options; using EventStore.Core.Data; +using EventStore.Core.Services.Transport.Grpc; using EventStore.Core.Tests; using EventStore.Core.Tests.Helpers; using EventStore.Core.Util; using EventStore.Projections.Core.Services.Processing; +using Google.Protobuf; +using Grpc.Core; +using Grpc.Net.Client; using NUnit.Framework; -using ExpectedVersion = EventStore.ClientAPI.ExpectedVersion; -using ResolvedEvent = EventStore.ClientAPI.ResolvedEvent; +using GrpcMetadata = EventStore.Core.Services.Transport.Grpc.Constants.Metadata; +using StatisticsReq = EventStore.Client.Projections.StatisticsReq; +using StreamsClient = EventStore.Client.Streams.Streams.StreamsClient; namespace EventStore.Projections.Core.Tests.ClientAPI.Cluster; -[Category("ClientAPI")] +[Category("Grpc")] public abstract class specification_with_standard_projections_runnning : SpecificationWithDirectoryPerTestFixture { + private static readonly TimeSpan PollTimeout = TimeSpan.FromSeconds(30); + private static readonly TimeSpan PollInterval = TimeSpan.FromMilliseconds(250); + private static readonly TimeSpan OperationTimeout = TimeSpan.FromMinutes(2); + private static readonly TimeSpan StandardProjectionsStartupTimeout = TimeSpan.FromMinutes(3); + private static readonly int PollAttemptCount = (int)(PollTimeout / PollInterval); protected MiniClusterNode[] _nodes = new MiniClusterNode[3]; protected Endpoints[] _nodeEndpoints = new Endpoints[3]; - protected IEventStoreConnection _conn; private readonly ProjectionsSubsystem[] _projections = new ProjectionsSubsystem[3]; - protected UserCredentials _admin = DefaultData.AdminCredentials; + private HttpClient _streamHttpClient; + private GrpcChannel _streamChannel; + private StreamsClient _streams; + private CancellationToken _operationCancellationToken; private protected ProjectionManagementTestClient ProjectionClient; protected class Endpoints @@ -39,13 +52,11 @@ protected class Endpoints public Endpoints(int internalTcp, int externalTcp, int httpPort) { var testIp = Environment.GetEnvironmentVariable("ES-TESTIP"); - var address = string.IsNullOrEmpty(testIp) ? IPAddress.Loopback : IPAddress.Parse(testIp); InternalTcp = new IPEndPoint(address, internalTcp); ExternalTcp = new IPEndPoint(address, externalTcp); HttpEndPoint = new IPEndPoint(address, httpPort); - - _ports = new[] { internalTcp, httpPort, externalTcp }; + _ports = [internalTcp, httpPort, externalTcp]; } public IEnumerable Ports => _ports; @@ -55,29 +66,23 @@ public Endpoints(int internalTcp, int externalTcp, int httpPort) public override async Task TestFixtureSetUp() { await base.TestFixtureSetUp(); -#if (!DEBUG) - Assert.Ignore("These tests require DEBUG conditional"); -#else - _nodeEndpoints[0] = new Endpoints( - PortsHelper.GetAvailablePort(IPAddress.Loopback), - PortsHelper.GetAvailablePort(IPAddress.Loopback), - PortsHelper.GetAvailablePort(IPAddress.Loopback)); - _nodeEndpoints[1] = new Endpoints( - PortsHelper.GetAvailablePort(IPAddress.Loopback), - PortsHelper.GetAvailablePort(IPAddress.Loopback), - PortsHelper.GetAvailablePort(IPAddress.Loopback)); - _nodeEndpoints[2] = new Endpoints( - PortsHelper.GetAvailablePort(IPAddress.Loopback), - PortsHelper.GetAvailablePort(IPAddress.Loopback), - PortsHelper.GetAvailablePort(IPAddress.Loopback)); - - _nodes[0] = CreateNode(0, - _nodeEndpoints[0], new[] { _nodeEndpoints[0].HttpEndPoint }); - _nodes[1] = CreateNode(1, - _nodeEndpoints[1], new[] { _nodeEndpoints[1].HttpEndPoint }); - _nodes[2] = CreateNode(2, - _nodeEndpoints[2], new[] { _nodeEndpoints[2].HttpEndPoint }); - WaitIdle(); + + for (var index = 0; index < _nodeEndpoints.Length; index++) + { + _nodeEndpoints[index] = new Endpoints( + PortsHelper.GetAvailablePort(IPAddress.Loopback), + PortsHelper.GetAvailablePort(IPAddress.Loopback), + PortsHelper.GetAvailablePort(IPAddress.Loopback)); + } + + for (var index = 0; index < _nodes.Length; index++) + { + var gossipSeeds = _nodeEndpoints + .Where((_, otherIndex) => otherIndex != index) + .Select(x => (EndPoint)x.HttpEndPoint) + .ToArray(); + _nodes[index] = CreateNode(index, _nodeEndpoints[index], gossipSeeds); + } var projectionsStarted = _projections.Select(p => SystemProjections.Created(p.LeaderInputBus)).ToArray(); @@ -87,50 +92,49 @@ public override async Task TestFixtureSetUp() node.WaitIdle(); } - await Task.WhenAll(_nodes.Select(x => x.Started)).WithTimeout(TimeSpan.FromSeconds(30)); - - _conn = EventStoreConnection.Create(_nodes[0].ExternalTcpEndPoint); - await _conn.ConnectAsync().WithTimeout(); + await Task.WhenAll(_nodes.Select(x => x.Started)).WithTimeout(TimeSpan.FromMinutes(5)); + await Task.WhenAll(_nodes.Select(x => x.AdminUserCreated)).WithTimeout(TimeSpan.FromMinutes(5)); var leader = _nodes.Single(x => x.NodeState == VNodeState.Leader); + _streamHttpClient = leader.CreateHttpClient(); + _streamChannel = GrpcChannel.ForAddress( + _streamHttpClient.BaseAddress, + new GrpcChannelOptions + { + HttpClient = _streamHttpClient, + DisposeHttpClient = false + }); + _streams = new StreamsClient(_streamChannel); ProjectionClient = new ProjectionManagementTestClient(leader.HttpEndPoint, leader.CreateHttpClient()); if (GivenStandardProjectionsRunning()) { - await Task.WhenAny(projectionsStarted).WithTimeout(TimeSpan.FromSeconds(10)); - await EnableStandardProjections().WithTimeout(TimeSpan.FromMinutes(2)); - } - - WaitIdle(); - - try - { - await Given().WithTimeout(); - } - catch (Exception ex) - { - throw new Exception("Given Failed", ex); + await Task.WhenAny(projectionsStarted).WithTimeout(OperationTimeout); + await RunBoundedOperation(EnableStandardProjections, StandardProjectionsStartupTimeout); } - try - { - await When().WithTimeout(); - } - catch (Exception ex) - { - throw new Exception("When Failed", ex); - } -#endif + await RunBoundedOperation(Given); + await RunBoundedOperation(When); } private MiniClusterNode CreateNode(int index, Endpoints endpoints, EndPoint[] gossipSeeds) { - _projections[index] = new ProjectionsSubsystem(new ProjectionSubsystemOptions(1, ProjectionType.All, false, TimeSpan.FromMinutes(Opts.ProjectionsQueryExpiryDefault), Opts.FaultOutOfOrderProjectionsDefault, 500, 250)); - var node = new MiniClusterNode( - PathName, index, endpoints.InternalTcp, - endpoints.ExternalTcp, endpoints.HttpEndPoint, - subsystems: [_projections[index]], gossipSeeds: gossipSeeds); - return node; + _projections[index] = new ProjectionsSubsystem(new ProjectionSubsystemOptions( + 1, + ProjectionType.All, + false, + TimeSpan.FromMinutes(Opts.ProjectionsQueryExpiryDefault), + Opts.FaultOutOfOrderProjectionsDefault, + 500, + 250)); + return new MiniClusterNode( + PathName, + index, + endpoints.InternalTcp, + endpoints.ExternalTcp, + endpoints.HttpEndPoint, + subsystems: [_projections[index]], + gossipSeeds: gossipSeeds); } [TearDown] @@ -159,52 +163,47 @@ protected async Task DisableStandardProjections() await DisableProjection(ProjectionNamesBuilder.StandardProjections.StreamsStandardProjection); } - protected virtual bool GivenStandardProjectionsRunning() - { - return true; - } + protected virtual bool GivenStandardProjectionsRunning() => true; protected async Task EnableProjection(string name) { - for (int i = 1; i <= 10; i++) + for (var attempt = 1; attempt <= 10; attempt++) { try { - await ProjectionClient.Enable(name); + await ProjectionClient.Enable(name, _operationCancellationToken); + break; } - catch (Exception) + catch when (attempt < 10 && !_operationCancellationToken.IsCancellationRequested) { - if (i == 10) - { - throw; - } - - await Task.Delay(5000); + await Task.Delay(500, _operationCancellationToken); } } - await Task.Delay(1000); /* workaround for race condition when multiple projections are being enabled simultaneously */ + await WaitForProjectionStatus(name, status => status.Contains("Running", StringComparison.OrdinalIgnoreCase)); } - protected Task DisableProjection(string name) + protected async Task DisableProjection(string name) { - return ProjectionClient.Disable(name); + await ProjectionClient.Disable(name, cancellationToken: _operationCancellationToken); + await WaitForProjectionStatus(name, status => status.Contains("Stopped", StringComparison.OrdinalIgnoreCase)); } - protected Task AbortProjection(string name) + protected async Task AbortProjection(string name) { - return ProjectionClient.Abort(name); + await ProjectionClient.Abort(name, _operationCancellationToken); + await WaitForProjectionStatus(name, status => status.StartsWith("Aborted", StringComparison.OrdinalIgnoreCase)); } [OneTimeTearDown] public override async Task TestFixtureTearDown() { ProjectionClient?.Dispose(); - _conn.Close(); - await Task.WhenAll( - _nodes[0].Shutdown(), - _nodes[1].Shutdown(), - _nodes[2].Shutdown()); + _streamChannel?.Dispose(); + _streamHttpClient?.Dispose(); + + await Task.WhenAll(_nodes.Where(x => x != null).Select(x => x.Shutdown())); + await base.TestFixtureTearDown(); } @@ -212,135 +211,216 @@ await Task.WhenAll( protected virtual Task Given() => Task.CompletedTask; - protected Task PostEvent(string stream, string eventType, string data) + protected async Task PostEvent(string stream, string eventType, string data) { - return _conn.AppendToStreamAsync(stream, ExpectedVersion.Any, new[] { CreateEvent(eventType, data) }); - } - - protected Task HardDeleteStream(string stream) - { - return _conn.DeleteStreamAsync(stream, ExpectedVersion.Any, true, _admin); + using var call = _streams.Append(GetCallOptions(_operationCancellationToken)); + await call.RequestStream.WriteAsync(new AppendReq + { + Options = new AppendReq.Types.Options + { + Any = new Empty(), + StreamIdentifier = new StreamIdentifier + { + StreamName = ByteString.CopyFromUtf8(stream) + } + } + }); + await call.RequestStream.WriteAsync(new AppendReq + { + ProposedMessage = new AppendReq.Types.ProposedMessage + { + Id = Uuid.NewUuid().ToDto(), + Data = ByteString.CopyFromUtf8(data), + CustomMetadata = ByteString.Empty, + Metadata = + { + { GrpcMetadata.Type, eventType }, + { GrpcMetadata.ContentType, GrpcMetadata.ContentTypes.ApplicationJson } + } + } + }); + await call.RequestStream.CompleteAsync(); + await call.ResponseAsync; } - protected Task SoftDeleteStream(string stream) + protected async Task HardDeleteStream(string stream) { - return _conn.DeleteStreamAsync(stream, ExpectedVersion.Any, false, _admin); + await _streams.TombstoneAsync(new TombstoneReq + { + Options = new TombstoneReq.Types.Options + { + Any = new Empty(), + StreamIdentifier = new StreamIdentifier + { + StreamName = ByteString.CopyFromUtf8(stream) + } + } + }, GetCallOptions(_operationCancellationToken)); } - protected static EventData CreateEvent(string type, string data) + protected async Task SoftDeleteStream(string stream) { - return new EventData(Guid.NewGuid(), type, true, Encoding.UTF8.GetBytes(data), Array.Empty()); + await _streams.DeleteAsync(new DeleteReq + { + Options = new DeleteReq.Types.Options + { + Any = new Empty(), + StreamIdentifier = new StreamIdentifier + { + StreamName = ByteString.CopyFromUtf8(stream) + } + } + }, GetCallOptions(_operationCancellationToken)); } protected void WaitIdle() { -#if DEBUG - _nodes[0].WaitIdle(); - _nodes[1].WaitIdle(); - _nodes[2].WaitIdle(); -#endif + foreach (var node in _nodes) + { + node.WaitIdle(); + } + + Thread.Sleep(50); } -#pragma warning disable 1998 protected async Task AssertStreamTailAsync(string streamId, params string[] events) { -#pragma warning restore 1998 -#if DEBUG - var result = await _conn.ReadStreamEventsBackwardAsync(streamId, -1, events.Length, true, _admin); - switch (result.Status) + string[] actual = []; + for (var attempt = 0; attempt < PollAttemptCount; attempt++) { - case SliceReadStatus.StreamDeleted: - Assert.Fail("Stream '{0}' is deleted", streamId); - break; - case SliceReadStatus.StreamNotFound: - Assert.Fail("Stream '{0}' does not exist", streamId); - break; - case SliceReadStatus.Success: - var resultEventsReversed = result.Events.Reverse().ToArray(); - if (resultEventsReversed.Length < events.Length) - { - DumpFailed("Stream does not contain enough events", streamId, events, result.Events); - } - else - { - for (var index = 0; index < events.Length; index++) - { - var parts = events[index].Split(new char[] { ':' }, 2); - var eventType = parts[0]; - var eventData = parts[1]; - - if (resultEventsReversed[index].Event.EventType != eventType) - { - DumpFailed("Invalid event type", streamId, events, resultEventsReversed); - } - else if (resultEventsReversed[index].Event.DebugDataView() != eventData) - { - DumpFailed("Invalid event body", streamId, events, resultEventsReversed); - } - } - } + actual = (await ReadStreamBackwards(streamId, (ulong)events.Length)) + .Reverse() + .Select(x => $"{x.Event.EventType()}:{x.Event.DebugDataView()}") + .ToArray(); + if (actual.SequenceEqual(events)) + { + return; + } - break; + await Task.Delay(PollInterval); } -#endif + + Assert.Fail( + $"Stream '{streamId}' did not reach the expected tail. Expected: [{string.Join(", ", events)}]. Actual: [{string.Join(", ", actual)}]."); } -#pragma warning disable 1998 protected async Task DumpStreamAsync(string streamId) { -#pragma warning restore 1998 -#if DEBUG - var result = await _conn.ReadStreamEventsBackwardAsync(streamId, -1, 100, true, _admin); - switch (result.Status) - { - case SliceReadStatus.StreamDeleted: - Assert.Fail("Stream '{0}' is deleted", streamId); - break; - case SliceReadStatus.StreamNotFound: - Assert.Fail("Stream '{0}' does not exist", streamId); - break; - case SliceReadStatus.Success: - Dump("Dumping..", streamId, result.Events.Reverse().ToArray()); - break; - } -#endif + var events = await ReadStreamBackwards(streamId, 100); + TestContext.Progress.WriteLine( + $"Stream '{streamId}': {string.Join(", ", events.Reverse().Select(x => $"{x.Event.EventType()}:{x.Event.DebugDataView()}"))}"); } -#if DEBUG - private void DumpFailed(string message, string streamId, string[] events, ResolvedEvent[] resultEvents) + protected async Task PostProjection(string query) { - var expected = events.Aggregate("", (a, v) => a + ", " + v); - var actual = resultEvents.Aggregate( - "", (a, v) => a + ", " + v.Event.EventType + ":" + v.Event.DebugDataView()); + await ProjectionClient.CreateContinuous( + "test-projection", + query, + cancellationToken: _operationCancellationToken); + await WaitForProjectionStatus( + "test-projection", + status => status.Contains("Running", StringComparison.OrdinalIgnoreCase)); + } - var actualMeta = resultEvents.Aggregate( - "", (a, v) => a + "\r\n" + v.Event.EventType + ":" + v.Event.DebugMetadataView()); + private async Task> ReadStreamBackwards(string streamId, ulong count) + { + using var call = _streams.Read(new ReadReq + { + Options = new ReadReq.Types.Options + { + Stream = new ReadReq.Types.Options.Types.StreamOptions + { + StreamIdentifier = new StreamIdentifier + { + StreamName = ByteString.CopyFromUtf8(streamId) + }, + End = new Empty() + }, + ReadDirection = ReadReq.Types.Options.Types.ReadDirection.Backwards, + ResolveLinks = true, + Count = count, + NoFilter = new Empty(), + UuidOption = new ReadReq.Types.Options.Types.UUIDOption + { + Structured = new Empty() + }, + ControlOption = new ReadReq.Types.Options.Types.ControlOption + { + Compatibility = 21 + } + } + }, GetCallOptions(_operationCancellationToken)); + var events = new List(); + while (await call.ResponseStream.MoveNext(_operationCancellationToken)) + { + if (call.ResponseStream.Current.ContentCase == ReadResp.ContentOneofCase.Event) + { + events.Add(call.ResponseStream.Current.Event); + } + } - Assert.Fail( - "Stream: '{0}'\r\n{1}\r\n\r\nExisting events: \r\n{2}\r\n Expected events: \r\n{3}\r\n\r\nActual metas:{4}", - streamId, - message, actual, expected, actualMeta); + return events; } - private void Dump(string message, string streamId, ResolvedEvent[] resultEvents) + private async Task WaitForProjectionStatus(string name, Func predicate) { - var actual = resultEvents.Aggregate( - "", (a, v) => a + ", " + v.OriginalEvent.EventType + ":" + v.OriginalEvent.DebugDataView()); + using var cancellation = CancellationTokenSource.CreateLinkedTokenSource(_operationCancellationToken); + cancellation.CancelAfter(PollTimeout); + var cancellationToken = cancellation.Token; + string lastStatus = null; + try + { + for (var attempt = 0; attempt < PollAttemptCount; attempt++) + { + var statistics = await ProjectionClient.Statistics(new StatisticsReq.Types.Options + { + Name = name + }, cancellationToken); + lastStatus = statistics.SingleOrDefault()?.Status; + if (lastStatus != null && predicate(lastStatus)) + { + return; + } - var actualMeta = resultEvents.Aggregate( - "", (a, v) => a + "\r\n" + v.OriginalEvent.EventType + ":" + v.OriginalEvent.DebugMetadataView()); + await Task.Delay(PollInterval, cancellationToken); + } + } + catch (OperationCanceledException) when (!_operationCancellationToken.IsCancellationRequested) + { + Assert.Fail($"Projection '{name}' did not reach the expected status. Last status: '{lastStatus}'."); + } - Debug.WriteLine( - "Stream: '{0}'\r\n{1}\r\n\r\nExisting events: \r\n{2}\r\n \r\nActual metas:{3}", streamId, - message, actual, actualMeta); + Assert.Fail($"Projection '{name}' did not reach the expected status. Last status: '{lastStatus}'."); } -#endif - protected async Task PostProjection(string query) + private async Task RunBoundedOperation(Func operation, TimeSpan? timeout = null) { - await ProjectionClient.CreateContinuous("test-projection", query); - WaitIdle(); + using var cancellation = new CancellationTokenSource(timeout ?? OperationTimeout); + _operationCancellationToken = cancellation.Token; + try + { + await operation().WaitAsync(cancellation.Token); + } + finally + { + _operationCancellationToken = default; + } + } + + private static CallOptions GetCallOptions(CancellationToken cancellationToken = default) + { + var credentials = CallCredentials.FromInterceptor((_, metadata) => + { + metadata.Add("authorization", + $"Basic {Convert.ToBase64String(Encoding.ASCII.GetBytes("admin:changeit"))}"); + return Task.CompletedTask; + }); + + return new CallOptions( + credentials: credentials, + deadline: DateTime.UtcNow.Add(PollTimeout), + cancellationToken: cancellationToken); } } @@ -351,7 +431,7 @@ public class vnode_cluster_specification : specification_ [Test, Explicit] public async Task vnode_cluster_starts() { - await PostProjection(@"fromStream('$user-admin').outputState()"); + await PostProjection(@"fromStream('$user-admin').when({$any:function(){return {}}}).outputState()"); await AssertStreamTailAsync("$projections-test-projection-result", "Result:{}"); } } diff --git a/src/EventStore.Projections.Core.Tests/ClientAPI/RecordedEventExtensions.cs b/src/EventStore.Projections.Core.Tests/ClientAPI/RecordedEventExtensions.cs index ae26152b90..b61c602376 100644 --- a/src/EventStore.Projections.Core.Tests/ClientAPI/RecordedEventExtensions.cs +++ b/src/EventStore.Projections.Core.Tests/ClientAPI/RecordedEventExtensions.cs @@ -1,10 +1,18 @@ -using System.Text; -using EventStore.ClientAPI; +using EventStore.Client.Streams; +using Google.Protobuf; +using GrpcMetadata = EventStore.Core.Services.Transport.Grpc.Constants.Metadata; +using RecordedEvent = EventStore.Client.Streams.ReadResp.Types.ReadEvent.Types.RecordedEvent; namespace EventStore.Projections.Core.Tests.ClientAPI; internal static class RecordedEventExtensions { - public static string DebugDataView(this RecordedEvent source) => Encoding.UTF8.GetString(source.Data); - public static string DebugMetadataView(this RecordedEvent source) => Encoding.UTF8.GetString(source.Metadata); + public static string DebugDataView(this RecordedEvent source) => source.Data.ToStringUtf8(); + + public static string DebugMetadataView(this RecordedEvent source) => source.CustomMetadata.ToStringUtf8(); + + public static string EventType(this RecordedEvent source) => + source.Metadata.TryGetValue(GrpcMetadata.Type, out var eventType) + ? eventType + : string.Empty; } diff --git a/src/EventStore.Projections.Core.Tests/ClientAPI/specification_with_standard_projections_runnning.cs b/src/EventStore.Projections.Core.Tests/ClientAPI/specification_with_standard_projections_runnning.cs index 2bad847175..45d4dd2ca1 100644 --- a/src/EventStore.Projections.Core.Tests/ClientAPI/specification_with_standard_projections_runnning.cs +++ b/src/EventStore.Projections.Core.Tests/ClientAPI/specification_with_standard_projections_runnning.cs @@ -1,41 +1,63 @@ using System; -using System.Diagnostics; +using System.Collections.Generic; using System.Linq; using System.Text; +using System.Threading; using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.ClientAPI.SystemData; +using EventStore.Client; +using EventStore.Client.Streams; using EventStore.Common.Options; +using EventStore.Core.Services.Transport.Grpc; using EventStore.Core.Tests; -using EventStore.Core.Tests.ClientAPI.Helpers; using EventStore.Core.Tests.Helpers; using EventStore.Core.Util; using EventStore.Projections.Core.Services.Processing; +using Google.Protobuf; +using Grpc.Core; +using Grpc.Net.Client; using NUnit.Framework; -using ResolvedEvent = EventStore.ClientAPI.ResolvedEvent; +using GrpcMetadata = EventStore.Core.Services.Transport.Grpc.Constants.Metadata; +using StatisticsReq = EventStore.Client.Projections.StatisticsReq; +using StreamsClient = EventStore.Client.Streams.Streams.StreamsClient; namespace EventStore.Projections.Core.Tests.ClientAPI; -[Category("ClientAPI")] +[Category("Grpc")] public abstract class specification_with_standard_projections_runnning : SpecificationWithDirectoryPerTestFixture { - protected IEventStoreConnection _conn; - protected UserCredentials _admin = DefaultData.AdminCredentials; - private protected ProjectionManagementTestClient ProjectionClient; - protected virtual TimeSpan StartupTimeout => TimeSpan.FromMinutes(5); + protected sealed class StreamReadResult + { + public StreamReadResult(bool exists, IReadOnlyList events) + { + Exists = exists; + Events = events; + } + + public bool Exists { get; } + public IReadOnlyList Events { get; } + } + private static readonly TimeSpan PollTimeout = TimeSpan.FromSeconds(20); + private static readonly TimeSpan PollInterval = TimeSpan.FromMilliseconds(250); + private static readonly TimeSpan OperationTimeout = TimeSpan.FromMinutes(2); + private static readonly int PollAttemptCount = (int)(PollTimeout / PollInterval); + private GrpcChannel _streamChannel; + private StreamsClient _streams; private Task _projectionsCreated; private ProjectionsSubsystem _projections; private MiniNode _node; + private CancellationToken _operationCancellationToken; + + private protected ProjectionManagementTestClient ProjectionClient; + protected virtual TimeSpan StartupTimeout => TimeSpan.FromMinutes(5); [OneTimeSetUp] public override async Task TestFixtureSetUp() { await base.TestFixtureSetUp(); - var projectionWorkerThreadCount = GivenWorkerThreadCount(); var configuration = new ProjectionSubsystemOptions( - projectionWorkerThreadCount, + GivenWorkerThreadCount(), ProjectionType.All, false, TimeSpan.FromMinutes(Opts.ProjectionsQueryExpiryDefault), @@ -49,25 +71,27 @@ public override async Task TestFixtureSetUp() _projectionsCreated = SystemProjections.Created(_projections.LeaderInputBus); await _node.Start(StartupTimeout); - await _node.WaitForTcpEndPoint().WithTimeout(StartupTimeout); - _conn = await TestConnectionLifecycle.ReconnectUntilReady( - CreateConnection, - connection => connection.ReadAllEventsForwardAsync(Position.Start, 1, false, _admin), - StartupTimeout); + await _node.AdminUserCreated.WithTimeout(StartupTimeout); + await _projectionsCreated.WithTimeout(PollTimeout); + _streamChannel = GrpcChannel.ForAddress( + _node.HttpClient.BaseAddress ?? new UriBuilder { Scheme = Uri.UriSchemeHttps }.Uri, + new GrpcChannelOptions + { + HttpClient = _node.HttpClient, + DisposeHttpClient = false + }); + _streams = new StreamsClient(_streamChannel); ProjectionClient = new ProjectionManagementTestClient(_node.HttpEndPoint, _node.HttpMessageHandler); - WaitIdle(); - if (GivenStandardProjectionsRunning()) { - await EnableStandardProjections(); + await RunBoundedOperation(EnableStandardProjections); } - WaitIdle(); try { - await Given().WithTimeout(TimeSpan.FromSeconds(10)); + await RunBoundedOperation(Given); } catch (Exception ex) { @@ -76,7 +100,7 @@ public override async Task TestFixtureSetUp() try { - await When().WithTimeout(TimeSpan.FromSeconds(10)); + await RunBoundedOperation(When); } catch (Exception ex) { @@ -84,10 +108,7 @@ public override async Task TestFixtureSetUp() } } - protected virtual int GivenWorkerThreadCount() - { - return 1; - } + protected virtual int GivenWorkerThreadCount() => 1; [TearDown] public async Task PostTestAsserts() @@ -101,7 +122,6 @@ public async Task PostTestAsserts() protected async Task EnableStandardProjections() { - await _projectionsCreated; await EnableProjection(ProjectionNamesBuilder.StandardProjections.EventByCategoryStandardProjection); await EnableProjection(ProjectionNamesBuilder.StandardProjections.EventByTypeStandardProjection); await EnableProjection(ProjectionNamesBuilder.StandardProjections.StreamByCategoryStandardProjection); @@ -116,64 +136,43 @@ protected async Task DisableStandardProjections() await DisableProjection(ProjectionNamesBuilder.StandardProjections.StreamsStandardProjection); } - protected virtual bool GivenStandardProjectionsRunning() - { - return true; - } + protected virtual bool GivenStandardProjectionsRunning() => true; - protected Task EnableProjection(string name) + protected async Task EnableProjection(string name) { - return ProjectionClient.Enable(name); + await ProjectionClient.Enable(name, _operationCancellationToken); + await WaitForProjectionStatus(name, status => status.Contains("Running", StringComparison.OrdinalIgnoreCase)); } - protected Task DisableProjection(string name) + protected async Task DisableProjection(string name) { - return ProjectionClient.Disable(name); + await ProjectionClient.Disable(name, cancellationToken: _operationCancellationToken); + await WaitForProjectionStatus(name, status => status.Contains("Stopped", StringComparison.OrdinalIgnoreCase)); } - protected Task AbortProjection(string name) + protected async Task AbortProjection(string name) { - return ProjectionClient.Abort(name); + await ProjectionClient.Abort(name, _operationCancellationToken); + await WaitForProjectionStatus(name, status => status.StartsWith("Aborted", StringComparison.OrdinalIgnoreCase)); } - protected Task CreateContinuousProjection(string name, string query) + protected async Task CreateContinuousProjection(string name, string query) { - return ProjectionClient.CreateContinuous(name, query); - } - - protected Task CreateTransientProjection(string name, string query) - { - return ProjectionClient.CreateTransient(name, query); + await ProjectionClient.CreateContinuous(name, query, cancellationToken: _operationCancellationToken); + await WaitForProjectionStatus(name, status => status.Contains("Running", StringComparison.OrdinalIgnoreCase)); } [OneTimeTearDown] public override async Task TestFixtureTearDown() { - if (_conn != null) - { - try - { - await TestConnectionLifecycle.CloseConnectionAndWait(_conn, TimeSpan.FromSeconds(20)); - } - catch - { - TestConnectionLifecycle.TryCloseConnection(_conn); - } - finally - { - TestConnectionLifecycle.DisposeIfNeeded(_conn); - } - } - ProjectionClient?.Dispose(); + _streamChannel?.Dispose(); if (_node != null) { await _node.Shutdown(); } - await Task.Delay(1000); - await base.TestFixtureTearDown(); } @@ -181,145 +180,277 @@ public override async Task TestFixtureTearDown() protected virtual Task Given() => Task.CompletedTask; - protected Task PostEvent(string stream, string eventType, string data) + protected async Task PostEvent(string stream, string eventType, string data) { - return _conn.AppendToStreamAsync(stream, ExpectedVersion.Any, CreateEvent(eventType, data)); + await Append(stream, eventType, data, new AppendReq.Types.Options { Any = new Empty() }); } - protected Task HardDeleteStream(string stream) + protected Task AppendToNewStream(string stream, string eventType, string data) => + Append(stream, eventType, data, new AppendReq.Types.Options { NoStream = new Empty() }); + + protected Task AppendToStream( + string stream, + ulong expectedRevision, + string eventType, + string data) => + Append(stream, eventType, data, new AppendReq.Types.Options { Revision = expectedRevision }); + + protected Task HardDeleteStream(string stream) => + Tombstone(stream, new TombstoneReq.Types.Options { Any = new Empty() }); + + protected Task HardDeleteStream(string stream, ulong expectedRevision) => + Tombstone(stream, new TombstoneReq.Types.Options { Revision = expectedRevision }); + + protected Task SoftDeleteStream(string stream) => + Delete(stream, new DeleteReq.Types.Options { Any = new Empty() }); + + protected Task SoftDeleteStream(string stream, ulong expectedRevision) => + Delete(stream, new DeleteReq.Types.Options { Revision = expectedRevision }); + + protected void WaitIdle(int multiplier = 1) { - return _conn.DeleteStreamAsync(stream, ExpectedVersion.Any, true, _admin); + _node.WaitIdle(); + Thread.Sleep(TimeSpan.FromMilliseconds(50 * multiplier)); } - protected Task SoftDeleteStream(string stream) + protected async Task AssertStreamTail(string streamId, params string[] events) { - return _conn.DeleteStreamAsync(stream, ExpectedVersion.Any, false, _admin); + string[] actual = []; + for (var attempt = 0; attempt < PollAttemptCount; attempt++) + { + var result = await ReadStream(streamId, (ulong)events.Length, true, true); + actual = result.Events + .Reverse() + .Select(FormatEvent) + .ToArray(); + + if (result.Exists && actual.SequenceEqual(events)) + { + return; + } + + await Task.Delay(PollInterval); + } + + Assert.Fail( + $"Stream '{streamId}' did not reach the expected tail. Expected: [{string.Join(", ", events)}]. Actual: [{string.Join(", ", actual)}]."); } - protected static EventData CreateEvent(string type, string data) + protected async Task ReadStreamForward(string streamId, ulong count, bool resolveLinks) => + await ReadStream(streamId, count, resolveLinks, false); + + protected async Task WaitForStreamEvents( + string streamId, + int minimumEventCount, + bool resolveLinks) { - return new EventData(Guid.NewGuid(), type, true, Encoding.UTF8.GetBytes(data), Array.Empty()); + StreamReadResult result = null; + for (var attempt = 0; attempt < PollAttemptCount; attempt++) + { + result = await ReadStreamForward(streamId, 100, resolveLinks); + if (result.Exists && result.Events.Count >= minimumEventCount) + { + return result; + } + + await Task.Delay(PollInterval); + } + + Assert.Fail( + $"Stream '{streamId}' did not contain {minimumEventCount} events. Actual: {result?.Events.Count ?? 0}."); + return result; } - private IEventStoreConnection CreateConnection() + protected async Task DumpStream(string streamId) { - return TestConnection.CreateMiniNodeClient(_node.TcpEndPoint); + var result = await ReadStreamForward(streamId, 100, true); + TestContext.Progress.WriteLine( + $"Stream '{streamId}': {string.Join(", ", result.Events.Select(FormatEvent))}"); } - protected void WaitIdle(int multiplier = 1) + protected async Task PostProjection(string query) { -#if DEBUG - _node.WaitIdle(); -#endif + await CreateContinuousProjection("test-projection", query); } -#pragma warning disable 1998 - protected async Task AssertStreamTail(string streamId, params string[] events) + private async Task Append( + string stream, + string eventType, + string data, + AppendReq.Types.Options options) { -#pragma warning restore 1998 -#if DEBUG - await Task.Delay(TimeSpan.FromMilliseconds(500)); - var result = await _conn.ReadStreamEventsBackwardAsync(streamId, -1, events.Length, true, _admin); - switch (result.Status) + options.StreamIdentifier = new StreamIdentifier { - case SliceReadStatus.StreamDeleted: - Assert.Fail("Stream '{0}' is deleted", streamId); - break; - case SliceReadStatus.StreamNotFound: - Assert.Fail("Stream '{0}' does not exist", streamId); - break; - case SliceReadStatus.Success: - var resultEventsReversed = result.Events.Reverse().ToArray(); - if (resultEventsReversed.Length < events.Length) - { - DumpFailed("Stream does not contain enough events", streamId, events, result.Events); - } - else + StreamName = ByteString.CopyFromUtf8(stream) + }; + + using var call = _streams.Append(GetCallOptions(_operationCancellationToken)); + await call.RequestStream.WriteAsync(new AppendReq { Options = options }); + await call.RequestStream.WriteAsync(new AppendReq + { + ProposedMessage = new AppendReq.Types.ProposedMessage + { + Id = Uuid.NewUuid().ToDto(), + Data = ByteString.CopyFromUtf8(data), + CustomMetadata = ByteString.Empty, + Metadata = { - for (var index = 0; index < events.Length; index++) - { - var parts = events[index].Split(new char[] { ':' }, 2); - var eventType = parts[0]; - var eventData = parts[1]; - - if (resultEventsReversed[index].Event.EventType != eventType) - { - DumpFailed("Invalid event type", streamId, events, resultEventsReversed); - } - else if (resultEventsReversed[index].Event.DebugDataView() != eventData) - { - DumpFailed("Invalid event body", streamId, events, resultEventsReversed); - } - } + { GrpcMetadata.Type, eventType }, + { GrpcMetadata.ContentType, GrpcMetadata.ContentTypes.ApplicationJson } } + } + }); + await call.RequestStream.CompleteAsync(); + return await call.ResponseAsync; + } - break; - } -#endif + private async Task Delete(string stream, DeleteReq.Types.Options options) + { + options.StreamIdentifier = new StreamIdentifier + { + StreamName = ByteString.CopyFromUtf8(stream) + }; + await _streams.DeleteAsync( + new DeleteReq { Options = options }, + GetCallOptions(_operationCancellationToken)); } -#pragma warning disable 1998 - protected async Task DumpStream(string streamId) + private async Task Tombstone(string stream, TombstoneReq.Types.Options options) { -#pragma warning restore 1998 -#if DEBUG - var result = await _conn.ReadStreamEventsBackwardAsync(streamId, -1, 100, true, _admin); - switch (result.Status) + options.StreamIdentifier = new StreamIdentifier { - case SliceReadStatus.StreamDeleted: - Assert.Fail("Stream '{0}' is deleted", streamId); - break; - case SliceReadStatus.StreamNotFound: - Assert.Fail("Stream '{0}' does not exist", streamId); - break; - case SliceReadStatus.Success: - Dump("Dumping..", streamId, result.Events.Reverse().ToArray()); - break; - } -#endif + StreamName = ByteString.CopyFromUtf8(stream) + }; + await _streams.TombstoneAsync( + new TombstoneReq { Options = options }, + GetCallOptions(_operationCancellationToken)); } -#if DEBUG - private void DumpFailed(string message, string streamId, string[] events, ResolvedEvent[] resultEvents) + private async Task ReadStream( + string streamId, + ulong count, + bool resolveLinks, + bool backwards) { - var expected = events.Aggregate("", (a, v) => a + ", " + v); - var actual = resultEvents.Aggregate( - "", (a, v) => a + ", " + v.Event.EventType + ":" + v.Event.DebugDataView()); + var stream = new ReadReq.Types.Options.Types.StreamOptions + { + StreamIdentifier = new StreamIdentifier + { + StreamName = ByteString.CopyFromUtf8(streamId) + } + }; + if (backwards) + { + stream.End = new Empty(); + } + else + { + stream.Start = new Empty(); + } - var actualMeta = resultEvents.Aggregate( - "", (a, v) => a + "\r\n" + v.Event.EventType + ":" + v.Event.DebugMetadataView()); + using var call = _streams.Read(new ReadReq + { + Options = new ReadReq.Types.Options + { + Stream = stream, + ReadDirection = backwards + ? ReadReq.Types.Options.Types.ReadDirection.Backwards + : ReadReq.Types.Options.Types.ReadDirection.Forwards, + ResolveLinks = resolveLinks, + Count = count, + NoFilter = new Empty(), + UuidOption = new ReadReq.Types.Options.Types.UUIDOption + { + Structured = new Empty() + }, + ControlOption = new ReadReq.Types.Options.Types.ControlOption + { + Compatibility = 21 + } + } + }, GetCallOptions(_operationCancellationToken)); + var exists = true; + var events = new List(); + while (await call.ResponseStream.MoveNext(_operationCancellationToken)) + { + switch (call.ResponseStream.Current.ContentCase) + { + case ReadResp.ContentOneofCase.Event: + events.Add(call.ResponseStream.Current.Event); + break; + case ReadResp.ContentOneofCase.StreamNotFound: + exists = false; + break; + } + } - Assert.Fail( - "Stream: '{0}'\r\n{1}\r\n\r\nExisting events: \r\n{2}\r\n Expected events: \r\n{3}\r\n\r\nActual metas:{4}", - streamId, - message, actual, expected, actualMeta); + return new StreamReadResult(exists, events); } - protected void Dump(string message, string streamId, ResolvedEvent[] resultEvents) + private async Task WaitForProjectionStatus(string name, Func predicate) { - var actual = resultEvents.Aggregate( - "", (a, v) => a + ", " + v.OriginalEvent.EventType + ":" + v.OriginalEvent.DebugDataView()); + using var cancellation = CancellationTokenSource.CreateLinkedTokenSource(_operationCancellationToken); + cancellation.CancelAfter(PollTimeout); + var cancellationToken = cancellation.Token; + string lastStatus = null; + try + { + for (var attempt = 0; attempt < PollAttemptCount; attempt++) + { + var statistics = await ProjectionClient.Statistics(new StatisticsReq.Types.Options + { + Name = name + }, cancellationToken); + lastStatus = statistics.SingleOrDefault()?.Status; + if (lastStatus != null && predicate(lastStatus)) + { + return; + } - var actualMeta = resultEvents.Aggregate( - "", (a, v) => a + "\r\n" + v.OriginalEvent.EventType + ":" + v.OriginalEvent.DebugMetadataView()); + await Task.Delay(PollInterval, cancellationToken); + } + } + catch (OperationCanceledException) when (!_operationCancellationToken.IsCancellationRequested) + { + Assert.Fail($"Projection '{name}' did not reach the expected status. Last status: '{lastStatus}'."); + } + Assert.Fail($"Projection '{name}' did not reach the expected status. Last status: '{lastStatus}'."); + } - Debug.WriteLine( - "Stream: '{0}'\r\n{1}\r\n\r\nExisting events: \r\n{2}\r\n \r\nActual metas:{3}", streamId, - message, actual, actualMeta); + private async Task RunBoundedOperation(Func operation) + { + using var cancellation = new CancellationTokenSource(OperationTimeout); + _operationCancellationToken = cancellation.Token; + try + { + await operation().WaitAsync(cancellation.Token); + } + finally + { + _operationCancellationToken = default; + } } -#endif - protected async Task PostProjection(string query) + private static string FormatEvent(ReadResp.Types.ReadEvent readEvent) { - await CreateContinuousProjection("test-projection", query); - WaitIdle(); + var recordedEvent = readEvent.Event ?? readEvent.Link; + return $"{recordedEvent.EventType()}:{recordedEvent.DebugDataView()}"; } - protected async Task PostQuery(string query) + private static CallOptions GetCallOptions(CancellationToken cancellationToken = default) { - await CreateTransientProjection("query", query); - WaitIdle(); + var credentials = CallCredentials.FromInterceptor((_, metadata) => + { + metadata.Add("authorization", + $"Basic {Convert.ToBase64String(Encoding.ASCII.GetBytes("admin:changeit"))}"); + return Task.CompletedTask; + }); + + return new CallOptions( + credentials: credentials, + deadline: DateTime.UtcNow.Add(PollTimeout), + cancellationToken: cancellationToken); } } diff --git a/src/EventStore.Projections.Core.Tests/ClientAPI/with_standard_projections_running.cs b/src/EventStore.Projections.Core.Tests/ClientAPI/with_standard_projections_running.cs index 6027793c96..dbec25fa71 100644 --- a/src/EventStore.Projections.Core.Tests/ClientAPI/with_standard_projections_running.cs +++ b/src/EventStore.Projections.Core.Tests/ClientAPI/with_standard_projections_running.cs @@ -1,10 +1,6 @@ -using System; using System.Text; using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.Core.Bus; using EventStore.Core.Tests; -using EventStore.Projections.Core.Services.Processing; using EventStore.Projections.Core.Services.Processing.Checkpointing; using Newtonsoft.Json.Linq; using NUnit.Framework; @@ -19,57 +15,59 @@ public abstract class when_deleting_stream_base [Test, Category("Network")] public async Task streams_stream_exists() { - Assert.AreEqual( - SliceReadStatus.Success, - (await _conn.ReadStreamEventsForwardAsync("$streams", 0, 10, false, _admin)).Status); + var result = await WaitForStreamEvents("$streams", 1, false); + Assert.That(result.Exists, Is.True); } [Test, Category("Network")] public async Task deleted_stream_events_are_indexed() { - await Task.Delay(500); //give the projection time to catchup... - var slice = await _conn.ReadStreamEventsForwardAsync("$ce-cat", 0, 10, true, _admin); - Assert.AreEqual(SliceReadStatus.Success, slice.Status); - - Assert.AreEqual(3, slice.Events.Length); - var deletedLinkMetadata = slice.Events[2].Link.Metadata; - Assert.IsNotNull(deletedLinkMetadata); - - var checkpointTag = Encoding.UTF8.GetString(deletedLinkMetadata).ParseCheckpointExtraJson(); - Assert.IsTrue(checkpointTag.TryGetValue("$deleted", out _)); - Assert.IsTrue(checkpointTag.TryGetValue("$o", out var originalStream)); - Assert.AreEqual("cat-1", ((JValue)originalStream).Value); + var result = await WaitForStreamEvents("$ce-cat", 3, true); + Assert.That(result.Events, Has.Count.EqualTo(3)); + + var deletedLink = result.Events[2].Link; + Assert.That(deletedLink, Is.Not.Null); + Assert.That(deletedLink.CustomMetadata, Is.Not.Null); + + var checkpointTag = Encoding.UTF8.GetString(deletedLink.CustomMetadata.ToByteArray()) + .ParseCheckpointExtraJson(); + Assert.That(checkpointTag.TryGetValue("$deleted", out _), Is.True); + Assert.That(checkpointTag.TryGetValue("$o", out var originalStream), Is.True); + Assert.That(((JValue)originalStream).Value, Is.EqualTo("cat-1")); } [Test, Category("Network")] public async Task deleted_stream_events_are_indexed_as_deleted() { - var slice = await _conn.ReadStreamEventsForwardAsync("$et-$deleted", 0, 10, true, _admin); - Assert.AreEqual(SliceReadStatus.Success, slice.Status); - - Assert.AreEqual(1, slice.Events.Length); + var result = await WaitForStreamEvents("$et-$deleted", 1, true); + Assert.That(result.Events, Has.Count.EqualTo(1)); } protected override async Task When() { await base.When(); - var r1 = await _conn.AppendToStreamAsync( - "cat-1", ExpectedVersion.NoStream, _admin, - new EventData(Guid.NewGuid(), "type1", true, Encoding.UTF8.GetBytes("{}"), null)) - ; - - var r2 = await _conn.AppendToStreamAsync( - "cat-1", r1.NextExpectedVersion, _admin, - new EventData(Guid.NewGuid(), "type1", true, Encoding.UTF8.GetBytes("{}"), null)); - - await _conn.DeleteStreamAsync("cat-1", r2.NextExpectedVersion, GivenDeleteHardDeleteStreamMode(), - _admin) - ; - WaitIdle(); + var firstAppend = await AppendToNewStream("cat-1", "type1", "{}"); + Assert.That(firstAppend.ResultCase, Is.EqualTo(EventStore.Client.Streams.AppendResp.ResultOneofCase.Success)); + + var secondAppend = await AppendToStream( + "cat-1", + firstAppend.Success.CurrentRevision, + "type1", + "{}"); + Assert.That(secondAppend.ResultCase, Is.EqualTo(EventStore.Client.Streams.AppendResp.ResultOneofCase.Success)); + + if (GivenDeleteHardDeleteStreamMode()) + { + await HardDeleteStream("cat-1", secondAppend.Success.CurrentRevision); + } + else + { + await SoftDeleteStream("cat-1", secondAppend.Success.CurrentRevision); + } + if (!GivenStandardProjectionsRunning()) { await EnableStandardProjections(); - WaitIdle(); } } @@ -79,47 +77,29 @@ await _conn.DeleteStreamAsync("cat-1", r2.NextExpectedVersion, GivenDeleteHardDe [TestFixture(typeof(LogFormat.V2), typeof(string))] public class when_hard_deleting_stream : when_deleting_stream_base { - protected override bool GivenDeleteHardDeleteStreamMode() - { - return true; - } + protected override bool GivenDeleteHardDeleteStreamMode() => true; } [TestFixture(typeof(LogFormat.V2), typeof(string))] public class when_soft_deleting_stream : when_deleting_stream_base { - protected override bool GivenDeleteHardDeleteStreamMode() - { - return false; - } + protected override bool GivenDeleteHardDeleteStreamMode() => false; } [TestFixture(typeof(LogFormat.V2), typeof(string))] public class when_hard_deleting_stream_and_starting_standard_projections : when_deleting_stream_base { - protected override bool GivenDeleteHardDeleteStreamMode() - { - return true; - } + protected override bool GivenDeleteHardDeleteStreamMode() => true; - protected override bool GivenStandardProjectionsRunning() - { - return false; - } + protected override bool GivenStandardProjectionsRunning() => false; } [TestFixture(typeof(LogFormat.V2), typeof(string))] public class when_soft_deleting_stream_and_starting_standard_projections : when_deleting_stream_base { - protected override bool GivenDeleteHardDeleteStreamMode() - { - return false; - } + protected override bool GivenDeleteHardDeleteStreamMode() => false; - protected override bool GivenStandardProjectionsRunning() - { - return false; - } + protected override bool GivenStandardProjectionsRunning() => false; } } } diff --git a/src/EventStore.Projections.Core.Tests/EventStore.Projections.Core.Tests.csproj b/src/EventStore.Projections.Core.Tests/EventStore.Projections.Core.Tests.csproj index 8cdd560ea4..93ac1fdb56 100644 --- a/src/EventStore.Projections.Core.Tests/EventStore.Projections.Core.Tests.csproj +++ b/src/EventStore.Projections.Core.Tests/EventStore.Projections.Core.Tests.csproj @@ -8,8 +8,6 @@ - - diff --git a/src/EventStore.Projections.Core.Tests/Playground/Launchpad.cs b/src/EventStore.Projections.Core.Tests/Playground/Launchpad.cs deleted file mode 100644 index 0441b45067..0000000000 --- a/src/EventStore.Projections.Core.Tests/Playground/Launchpad.cs +++ /dev/null @@ -1,86 +0,0 @@ -using System; -using System.Collections.Generic; -using System.IO; -using System.Runtime.InteropServices; -using System.Threading; -using EventStore.Common.Utils; -using NUnit.Framework; - -namespace EventStore.Projections.Core.Tests.Playground { - [TestFixture, Explicit, Category("Manual")] - public class Launchpad : LaunchpadBase { - private IDisposable _vnodeProcess; - private IDisposable _managerProcess; - private IDisposable _clientProcess; - private IDisposable _projectionsProcess; - - private string _binFolder; - private Dictionary _environment; - private string _dbPath; - - [SetUp] - public void Setup() { - if (!OS.IsUnix) - AllocConsole(); // this is required to keep console open after executeassemly has exited - - _binFolder = AppDomain.CurrentDomain.BaseDirectory; - _dbPath = Path.Combine(_binFolder, DateTime.UtcNow.Ticks.ToString()); - _environment = new Dictionary {{"EVENTSTORE_LOGSDIR", _dbPath}}; - - var vnodeExecutable = Path.Combine(_binFolder, @"EventStore.ClusterNode.exe"); - var managerExecutable = Path.Combine(_binFolder, @"EventStore.Manager.exe"); - - - string vnodeCommandLine = - string.Format( - @"--ip=127.0.0.1 --db={0} --int-tcp-port=3111 --ext-tcp-port=1111 --http-port=2111 --manager-ip=127.0.0.1 --manager-port=30777 --nodes-count=1 --fake-dns --prepare-count=1 --commit-count=1", - _dbPath); - string managerCommandLine = @"--ip=127.0.0.1 --port=30777 --nodes-count=1 --fake-dns"; - - _managerProcess = _launch(managerExecutable, managerCommandLine, _environment); - Thread.Sleep(500); - _vnodeProcess = _launch(vnodeExecutable, vnodeCommandLine, _environment); - } - - [TearDown] - public void Teardown() { - if (_managerProcess != null) _managerProcess.Dispose(); - if (_vnodeProcess != null) _vnodeProcess.Dispose(); - if (_clientProcess != null) _clientProcess.Dispose(); - if (_projectionsProcess != null) _projectionsProcess.Dispose(); - } - - public void LaunchFlood() { - var clientExecutable = Path.Combine(_binFolder, @"EventStore.Client.exe"); - string clientCommandLine = @"-i 127.0.0.1 -t 1111 -h 2111 WRFL 1 100"; - _clientProcess = _launch(clientExecutable, clientCommandLine, _environment); - } - - public void LaunchProjections() { - var clientExecutable = Path.Combine(_binFolder, @"EventStore.Projections.Worker.exe"); - string clientCommandLine = string.Format(@"--ip 127.0.0.1 -t 1111 -h 2111 --db {0}", _dbPath); - _projectionsProcess = _launch(clientExecutable, clientCommandLine, _environment); - } - - [Test] - public void WriteFloodAndProjections() { - Thread.Sleep(5000); - LaunchFlood(); - Thread.Sleep(5000); - LaunchProjections(); - Thread.Sleep(500); - LaunchFlood(); - Thread.Sleep(160000); - } - - [Test] - public void JustProjections() { - Thread.Sleep(3500); - LaunchProjections(); - Thread.Sleep(50000); - } - - [DllImport("kernel32.dll", SetLastError = true)] - private static extern bool AllocConsole(); - } -} diff --git a/src/EventStore.Projections.Core.Tests/Playground/Launchpad2.cs b/src/EventStore.Projections.Core.Tests/Playground/Launchpad2.cs deleted file mode 100644 index 0718dfeb09..0000000000 --- a/src/EventStore.Projections.Core.Tests/Playground/Launchpad2.cs +++ /dev/null @@ -1,68 +0,0 @@ -using System; -using System.Collections.Generic; -using System.IO; -using System.Runtime.InteropServices; -using System.Threading; -using EventStore.Common.Utils; -using NUnit.Framework; - -namespace EventStore.Projections.Core.Tests.Playground { - [TestFixture, Explicit, Category("Manual")] - public class Launchpad2 : LaunchpadBase { - private IDisposable _vnodeProcess; - private IDisposable _clientProcess; - - private string _binFolder; - private Dictionary _environment; - private string _dbPath; - - [SetUp] - public void Setup() { - if (!OS.IsUnix) - AllocConsole(); // this is required to keep console open after executeassemly has exited - - _binFolder = AppDomain.CurrentDomain.BaseDirectory; - _dbPath = Path.Combine(_binFolder, DateTime.UtcNow.Ticks.ToString()); - _environment = new Dictionary {{"EVENTSTORE_LOGSDIR", _dbPath}}; - - var vnodeExecutable = Path.Combine(_binFolder, @"EventStore.Projections.Worker.exe"); - - - string vnodeCommandLine = - string.Format( - @"--ip=127.0.0.1 --db={0} --stats-frequency-sec=10 --int-tcp-port=3111 --ext-tcp-port=1111 --http-port=2111", - _dbPath); - - _vnodeProcess = _launch(vnodeExecutable, vnodeCommandLine, _environment); - } - - [TearDown] - public void Teardown() { - if (_vnodeProcess != null) _vnodeProcess.Dispose(); - if (_clientProcess != null) _clientProcess.Dispose(); - } - - public void LaunchFlood() { - var clientExecutable = Path.Combine(_binFolder, @"EventStore.Client.exe"); - string clientCommandLine = @"--ip 127.0.0.1 --tcp-port 1111 --http-port 2111 WRFL 1 100"; - _clientProcess = _launch(clientExecutable, clientCommandLine, _environment); - } - - [Test] - public void RunSingle() { - Thread.Sleep(60000); - } - - [Test] - public void RunSingleAndFlood() { - Thread.Sleep(4000); - LaunchFlood(); - Thread.Sleep(10000); - LaunchFlood(); - Thread.Sleep(20000); - } - - [DllImport("kernel32.dll", SetLastError = true)] - private static extern bool AllocConsole(); - } -} diff --git a/src/EventStore.Projections.Core.Tests/ProjectionManagementTestClient.cs b/src/EventStore.Projections.Core.Tests/ProjectionManagementTestClient.cs index 18c7f9672c..3fbc0d5bce 100644 --- a/src/EventStore.Projections.Core.Tests/ProjectionManagementTestClient.cs +++ b/src/EventStore.Projections.Core.Tests/ProjectionManagementTestClient.cs @@ -42,7 +42,12 @@ public ProjectionManagementTestClient(IPEndPoint endpoint, HttpClient httpClient _client = new ProjectionGrpc.ProjectionsClient(_channel); } - public Task CreateContinuous(string name, string query, bool emitEnabled = false, bool trackEmittedStreams = false) => + public Task CreateContinuous( + string name, + string query, + bool emitEnabled = false, + bool trackEmittedStreams = false, + CancellationToken cancellationToken = default) => _client.CreateAsync(new CreateReq { Options = new CreateReq.Types.Options @@ -55,9 +60,9 @@ public Task CreateContinuous(string name, string query, bool emitEnabled = false }, Query = query } - }, GetCallOptions()).ResponseAsync; + }, GetCallOptions(cancellationToken)).ResponseAsync; - public Task CreateOneTime(string query, string name = null) => + public Task CreateOneTime(string query, string name = null, CancellationToken cancellationToken = default) => _client.CreateAsync(new CreateReq { Options = new CreateReq.Types.Options @@ -66,9 +71,9 @@ public Task CreateOneTime(string query, string name = null) => Name = name ?? string.Empty, Query = query } - }, GetCallOptions()).ResponseAsync; + }, GetCallOptions(cancellationToken)).ResponseAsync; - public Task CreateTransient(string name, string query) => + public Task CreateTransient(string name, string query, CancellationToken cancellationToken = default) => _client.CreateAsync(new CreateReq { Options = new CreateReq.Types.Options @@ -79,18 +84,21 @@ public Task CreateTransient(string name, string query) => }, Query = query } - }, GetCallOptions()).ResponseAsync; + }, GetCallOptions(cancellationToken)).ResponseAsync; - public Task Enable(string name) => + public Task Enable(string name, CancellationToken cancellationToken = default) => _client.EnableAsync(new EnableReq { Options = new EnableReq.Types.Options { Name = name } - }, GetCallOptions()).ResponseAsync; + }, GetCallOptions(cancellationToken)).ResponseAsync; - public Task Disable(string name, bool writeCheckpoint = true) => + public Task Disable( + string name, + bool writeCheckpoint = true, + CancellationToken cancellationToken = default) => _client.DisableAsync(new DisableReq { Options = new DisableReq.Types.Options @@ -98,18 +106,18 @@ public Task Disable(string name, bool writeCheckpoint = true) => Name = name, WriteCheckpoint = writeCheckpoint } - }, GetCallOptions()).ResponseAsync; + }, GetCallOptions(cancellationToken)).ResponseAsync; - public Task Abort(string name) => + public Task Abort(string name, CancellationToken cancellationToken = default) => _client.AbortAsync(new AbortReq { Options = new AbortReq.Types.Options { Name = name } - }, GetCallOptions()).ResponseAsync; + }, GetCallOptions(cancellationToken)).ResponseAsync; - public Task UpdateQuery(string name, string query) => + public Task UpdateQuery(string name, string query, CancellationToken cancellationToken = default) => _client.UpdateAsync(new UpdateReq { Options = new UpdateReq.Types.Options @@ -117,7 +125,7 @@ public Task UpdateQuery(string name, string query) => Name = name, Query = query } - }, GetCallOptions()).ResponseAsync; + }, GetCallOptions(cancellationToken)).ResponseAsync; public async Task> Statistics( StatisticsReq.Types.Options options, @@ -126,7 +134,7 @@ public Task UpdateQuery(string name, string query) => using var call = _client.Statistics(new StatisticsReq { Options = options - }, GetCallOptions(TimeSpan.FromSeconds(20))); + }, GetCallOptions(cancellationToken, TimeSpan.FromSeconds(20))); var results = new List(); while (await call.ResponseStream.MoveNext(cancellationToken)) @@ -161,7 +169,9 @@ public void Dispose() _ownedHttpClient?.Dispose(); } - private static CallOptions GetCallOptions(TimeSpan? deadline = null) + private static CallOptions GetCallOptions( + CancellationToken cancellationToken = default, + TimeSpan? deadline = null) { var credentials = CallCredentials.FromInterceptor((_, metadata) => { @@ -174,6 +184,7 @@ private static CallOptions GetCallOptions(TimeSpan? deadline = null) credentials: credentials, deadline: deadline is { } value ? DateTime.UtcNow.Add(value) - : null); + : null, + cancellationToken: cancellationToken); } } diff --git a/src/EventStore.Projections.Core.Tests/Services/SpecificationWithEmittedStreamsTrackerAndDeleter.cs b/src/EventStore.Projections.Core.Tests/Services/SpecificationWithEmittedStreamsTrackerAndDeleter.cs index 22ae22be66..dd37d86a39 100644 --- a/src/EventStore.Projections.Core.Tests/Services/SpecificationWithEmittedStreamsTrackerAndDeleter.cs +++ b/src/EventStore.Projections.Core.Tests/Services/SpecificationWithEmittedStreamsTrackerAndDeleter.cs @@ -1,15 +1,46 @@ +using System; +using System.Collections.Generic; +using System.IO; +using System.Linq; +using System.Text; +using System.Threading; using System.Threading.Tasks; +using EventStore.Client.Streams; using EventStore.Core.Helpers; using EventStore.Core.Messages; -using EventStore.Core.Tests.ClientAPI; +using EventStore.Core.Services.Transport.Grpc; +using EventStore.Core.Tests; +using EventStore.Core.Tests.Helpers; using EventStore.Projections.Core.Services; using EventStore.Projections.Core.Services.Processing; using EventStore.Projections.Core.Services.Processing.Emitting; +using Google.Protobuf; +using Grpc.Core; +using Grpc.Net.Client; +using NUnit.Framework; +using GrpcMetadata = EventStore.Core.Services.Transport.Grpc.Constants.Metadata; +using ReadEvent = EventStore.Client.Streams.ReadResp.Types.ReadEvent; +using StreamsClient = EventStore.Client.Streams.Streams.StreamsClient; namespace EventStore.Projections.Core.Tests.Services; -public abstract class SpecificationWithEmittedStreamsTrackerAndDeleter : SpecificationWithMiniNode +public abstract class SpecificationWithEmittedStreamsTrackerAndDeleter : SpecificationWithDirectoryPerTestFixture { + protected sealed class StreamReadResult + { + public StreamReadResult(bool exists, ReadEvent[] events) + { + Exists = exists; + Events = events; + } + + public bool Exists { get; } + public ReadEvent[] Events { get; } + } + + private GrpcChannel _channel; + protected MiniNode _node; + protected StreamsClient _client; protected IEmittedStreamsTracker _emittedStreamsTracker; protected IEmittedStreamsDeleter _emittedStreamsDeleter; protected ProjectionNamesBuilder _projectionNamesBuilder; @@ -17,8 +48,33 @@ public abstract class SpecificationWithEmittedStreamsTrackerAndDeleter(PathName); + await _node.Start(); + await _node.AdminUserCreated.WithTimeout(Timeout); + _channel = GrpcChannel.ForAddress(new UriBuilder { Scheme = Uri.UriSchemeHttps }.Uri, + new GrpcChannelOptions { HttpClient = _node.HttpClient, DisposeHttpClient = false }); + _client = new StreamsClient(_channel); + await Given().WithTimeout(Timeout); + await When().WithTimeout(Timeout); + } + + [OneTimeTearDown] + public override async Task TestFixtureTearDown() + { + _channel?.Dispose(); + await _node.Shutdown(); + await base.TestFixtureTearDown(); + } + + protected virtual Task Given() { _ioDispatcher = new IODispatcher(_node.Node.MainQueue, _node.Node.MainQueue, true); _node.Node.MainBus.Subscribe(_ioDispatcher.BackwardReader); @@ -38,4 +94,128 @@ protected override Task Given() _projectionNamesBuilder.GetEmittedStreamsCheckpointName()); return Task.CompletedTask; } + + protected async Task AppendEvent(string stream, string eventType, byte[] data) + { + using var call = _client.Append(AdminCallOptions()); + await call.RequestStream.WriteAsync(new AppendReq + { + Options = new() + { + Any = new(), + StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(stream) } + } + }); + await call.RequestStream.WriteAsync(new AppendReq + { + ProposedMessage = new() + { + Id = Uuid.NewUuid().ToDto(), + Data = ByteString.CopyFrom(data), + CustomMetadata = ByteString.Empty, + Metadata = { + { GrpcMetadata.Type, eventType }, + { GrpcMetadata.ContentType, GrpcMetadata.ContentTypes.ApplicationJson } + } + } + }); + await call.RequestStream.CompleteAsync(); + await call.ResponseAsync; + } + + protected async Task ReadEvents(string stream, int count) + { + using var timeout = new CancellationTokenSource(Timeout); + return await ReadEvents(stream, count, timeout.Token); + } + + protected AsyncServerStreamingCall SubscribeToStream( + string stream, + bool resolveLinks, + CancellationToken cancellationToken) => + _client.Read(new ReadReq + { + Options = new() + { + Stream = new() + { + StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(stream) }, + End = new() + }, + ReadDirection = ReadReq.Types.Options.Types.ReadDirection.Forwards, + ResolveLinks = resolveLinks, + Subscription = new(), + NoFilter = new(), + UuidOption = new() { Structured = new() } + } + }, AdminCallOptions(cancellationToken)); + + private async Task ReadEvents(string stream, int count, CancellationToken cancellationToken) + { + using var call = _client.Read(new ReadReq + { + Options = new() + { + Stream = new() + { + StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(stream) }, + Start = new() + }, + ReadDirection = ReadReq.Types.Options.Types.ReadDirection.Forwards, + Count = (ulong)count, + NoFilter = new(), + UuidOption = new() { Structured = new() } + } + }, AdminCallOptions(cancellationToken)); + var exists = true; + var events = new List(); + while (await call.ResponseStream.MoveNext(cancellationToken)) + { + switch (call.ResponseStream.Current.ContentCase) + { + case ReadResp.ContentOneofCase.Event: + events.Add(call.ResponseStream.Current.Event); + break; + case ReadResp.ContentOneofCase.StreamNotFound: + exists = false; + break; + } + } + + return new StreamReadResult(exists, events.ToArray()); + } + + protected async Task WaitForEvents(string stream, int count) + { + using var timeout = new CancellationTokenSource(Timeout); + var events = Array.Empty(); + try + { + while (true) + { + events = (await ReadEvents(stream, count, timeout.Token)).Events; + if (events.Length >= count) + return events; + await Task.Delay(50, timeout.Token); + } + } + catch (OperationCanceledException) when (timeout.IsCancellationRequested) + { + return events; + } + catch (RpcException ex) when (timeout.IsCancellationRequested && + ex.StatusCode is StatusCode.Cancelled or StatusCode.DeadlineExceeded) + { + return events; + } + } + + private static CallOptions AdminCallOptions(CancellationToken cancellationToken = default) => new( + credentials: CallCredentials.FromInterceptor((_, metadata) => + { + metadata.Add("authorization", + $"Basic {Convert.ToBase64String(Encoding.ASCII.GetBytes("admin:changeit"))}"); + return Task.CompletedTask; + }), + cancellationToken: cancellationToken); } diff --git a/src/EventStore.Projections.Core.Tests/Services/SpecificationWithEmittedStreamsTrackerAndDeleterTests.cs b/src/EventStore.Projections.Core.Tests/Services/SpecificationWithEmittedStreamsTrackerAndDeleterTests.cs new file mode 100644 index 0000000000..58e045c90b --- /dev/null +++ b/src/EventStore.Projections.Core.Tests/Services/SpecificationWithEmittedStreamsTrackerAndDeleterTests.cs @@ -0,0 +1,99 @@ +using System; +using System.Threading; +using System.Threading.Tasks; +using EventStore.Client.Streams; +using EventStore.Core.Tests; +using Grpc.Core; +using NUnit.Framework; +using StreamsClient = EventStore.Client.Streams.Streams.StreamsClient; + +namespace EventStore.Projections.Core.Tests.Services; + +[TestFixture] +public class SpecificationWithEmittedStreamsTrackerAndDeleterTests +{ + [Test] + public void read_events_cancels_the_stream_when_the_helper_times_out() + { + var client = new BlockingStreamsClient(); + var specification = new TestSpecification(client); + + Assert.ThrowsAsync(async () => + await specification.ReadEventCount("stream", 1)); + Assert.That(client.CallCancellationToken.CanBeCanceled, Is.True); + Assert.That(client.CallCancellationToken.IsCancellationRequested, Is.True); + Assert.That(client.Reader.MoveNextCancellationToken, + Is.EqualTo(client.CallCancellationToken)); + } + + [Test] + public async Task wait_for_events_returns_the_last_result_when_the_helper_times_out() + { + var client = new BlockingStreamsClient(); + var specification = new TestSpecification(client); + + var eventCount = await specification.WaitForEventCount("stream", 1); + + Assert.That(eventCount, Is.Zero); + Assert.That(client.CallCancellationToken.CanBeCanceled, Is.True); + Assert.That(client.CallCancellationToken.IsCancellationRequested, Is.True); + Assert.That(client.Reader.MoveNextCancellationToken, + Is.EqualTo(client.CallCancellationToken)); + } + + private sealed class TestSpecification + : SpecificationWithEmittedStreamsTrackerAndDeleter + { + public TestSpecification(StreamsClient client) + { + _client = client; + } + + protected override TimeSpan Timeout => TimeSpan.FromMilliseconds(25); + + protected override Task When() => Task.CompletedTask; + + public async Task ReadEventCount(string stream, int count) => + (await ReadEvents(stream, count)).Events.Length; + + public async Task WaitForEventCount(string stream, int count) => + (await WaitForEvents(stream, count)).Length; + } + + private sealed class BlockingStreamsClient : StreamsClient + { + public BlockingStreamReader Reader { get; } = new(); + public CancellationToken CallCancellationToken { get; private set; } + + public override AsyncServerStreamingCall Read(ReadReq request, CallOptions options) + { + CallCancellationToken = options.CancellationToken; + return new AsyncServerStreamingCall( + Reader, + Task.FromResult(new Metadata()), + () => Status.DefaultSuccess, + () => new Metadata(), + () => { }); + } + } + + private sealed class BlockingStreamReader : IAsyncStreamReader + { + public ReadResp Current => null; + public CancellationToken MoveNextCancellationToken { get; private set; } + + public Task MoveNext(CancellationToken cancellationToken) + { + MoveNextCancellationToken = cancellationToken; + return cancellationToken.CanBeCanceled + ? WaitForCancellation(cancellationToken) + : Task.FromResult(false); + } + + private static async Task WaitForCancellation(CancellationToken cancellationToken) + { + await Task.Delay(Timeout.InfiniteTimeSpan, cancellationToken); + return false; + } + } +} diff --git a/src/EventStore.Projections.Core.Tests/Services/emitted_streams_deleter/when_deleting/with_an_existing_emitted_streams_stream.cs b/src/EventStore.Projections.Core.Tests/Services/emitted_streams_deleter/when_deleting/with_an_existing_emitted_streams_stream.cs index 2db800c545..6aef0e4b11 100644 --- a/src/EventStore.Projections.Core.Tests/Services/emitted_streams_deleter/when_deleting/with_an_existing_emitted_streams_stream.cs +++ b/src/EventStore.Projections.Core.Tests/Services/emitted_streams_deleter/when_deleting/with_an_existing_emitted_streams_stream.cs @@ -1,7 +1,6 @@ using System; using System.Threading; using System.Threading.Tasks; -using EventStore.ClientAPI; using EventStore.Core.Tests; using EventStore.Projections.Core.Services.Processing; using EventStore.Projections.Core.Services.Processing.Checkpointing; @@ -18,19 +17,12 @@ public class with_an_existing_emitted_streams_stream : Sp protected ManualResetEvent _resetEvent = new ManualResetEvent(false); private string _testStreamName = "test_stream"; private ManualResetEvent _eventAppeared = new ManualResetEvent(false); - private EventStore.ClientAPI.SystemData.UserCredentials _credentials; protected override async Task Given() { - _credentials = new EventStore.ClientAPI.SystemData.UserCredentials("admin", "changeit"); _onDeleteStreamCompleted = () => { _resetEvent.Set(); }; await base.Given(); - var sub = await _conn.SubscribeToStreamAsync(_projectionNamesBuilder.GetEmittedStreamsName(), true, (s, evnt) => - { - _eventAppeared.Set(); - return Task.CompletedTask; - }, userCredentials: _credentials); _emittedStreamsTracker.TrackEmittedStream(new EmittedEvent[] { new EmittedDataEvent( @@ -38,18 +30,12 @@ protected override async Task Given() "data", null, CheckpointTag.FromPosition(0, 100, 50), null), }); - if (!_eventAppeared.WaitOne(TimeSpan.FromSeconds(5))) + var events = await WaitForEvents(_projectionNamesBuilder.GetEmittedStreamsName(), 1); + if (events.Length != 1) { Assert.Fail("Timed out waiting for emitted stream event"); } - - sub.Unsubscribe(); - - var emittedStreamResult = - await _conn.ReadStreamEventsForwardAsync(_projectionNamesBuilder.GetEmittedStreamsName(), 0, 1, false, - _credentials); - Assert.AreEqual(1, emittedStreamResult.Events.Length); - Assert.AreEqual(SliceReadStatus.Success, emittedStreamResult.Status); + _eventAppeared.Set(); } protected override Task When() @@ -66,25 +52,22 @@ protected override Task When() [Test] public async Task should_have_deleted_the_tracked_emitted_stream() { - var result = await _conn.ReadStreamEventsForwardAsync(_testStreamName, 0, 1, false, - new EventStore.ClientAPI.SystemData.UserCredentials("admin", "changeit")); - Assert.AreEqual(SliceReadStatus.StreamNotFound, result.Status); + var result = await ReadEvents(_testStreamName, 1); + Assert.That(result.Exists, Is.False); } [Test] public async Task should_have_deleted_the_checkpoint_stream() { - var result = await _conn.ReadStreamEventsForwardAsync(_projectionNamesBuilder.GetEmittedStreamsCheckpointName(), - 0, 1, false, new EventStore.ClientAPI.SystemData.UserCredentials("admin", "changeit")); - Assert.AreEqual(SliceReadStatus.StreamNotFound, result.Status); + var result = await ReadEvents(_projectionNamesBuilder.GetEmittedStreamsCheckpointName(), 1); + Assert.That(result.Exists, Is.False); } [Test] public async Task should_have_deleted_the_emitted_streams_stream() { - var result = await _conn.ReadStreamEventsForwardAsync(_projectionNamesBuilder.GetEmittedStreamsName(), 0, 1, - false, new EventStore.ClientAPI.SystemData.UserCredentials("admin", "changeit")); - Assert.AreEqual(SliceReadStatus.StreamNotFound, result.Status); + var result = await ReadEvents(_projectionNamesBuilder.GetEmittedStreamsName(), 1); + Assert.That(result.Exists, Is.False); } } diff --git a/src/EventStore.Projections.Core.Tests/Services/emitted_streams_deleter/when_deleting/with_multiple_tracked_streams.cs b/src/EventStore.Projections.Core.Tests/Services/emitted_streams_deleter/when_deleting/with_multiple_tracked_streams.cs index 68fbff7581..691bf2eb62 100644 --- a/src/EventStore.Projections.Core.Tests/Services/emitted_streams_deleter/when_deleting/with_multiple_tracked_streams.cs +++ b/src/EventStore.Projections.Core.Tests/Services/emitted_streams_deleter/when_deleting/with_multiple_tracked_streams.cs @@ -1,7 +1,6 @@ using System; using System.Threading; using System.Threading.Tasks; -using EventStore.ClientAPI; using EventStore.Common.Utils; using EventStore.Core.Tests; using EventStore.Projections.Core.Services.Processing; @@ -20,25 +19,16 @@ public class with_multiple_tracked_streams : Specificatio protected CountdownEvent _eventAppeared; private int _numberOfTrackedEvents = 50; private string _testStreamFormat = "test_stream_{0}"; - private EventStore.ClientAPI.SystemData.UserCredentials _credentials; protected override async Task Given() { - _credentials = new EventStore.ClientAPI.SystemData.UserCredentials("admin", "changeit"); _eventAppeared = new CountdownEvent(_numberOfTrackedEvents); _onDeleteStreamCompleted = () => { _resetEvent.Set(); }; await base.Given(); - var sub = await _conn.SubscribeToStreamAsync(_projectionNamesBuilder.GetEmittedStreamsName(), true, (s, evnt) => - { - _eventAppeared.Signal(); - return Task.CompletedTask; - }, userCredentials: _credentials); - for (int i = 0; i < _numberOfTrackedEvents; i++) { - await _conn.AppendToStreamAsync(String.Format(_testStreamFormat, i), ExpectedVersion.Any, - new EventData(Guid.NewGuid(), "type1", true, Helper.UTF8NoBom.GetBytes("data"), null)); + await AppendEvent(String.Format(_testStreamFormat, i), "type1", Helper.UTF8NoBom.GetBytes("data")); _emittedStreamsTracker.TrackEmittedStream(new EmittedEvent[] { new EmittedDataEvent( String.Format(_testStreamFormat, i), Guid.NewGuid(), "type1", true, @@ -46,16 +36,13 @@ await _conn.AppendToStreamAsync(String.Format(_testStreamFormat, i), ExpectedVer }); } - if (!_eventAppeared.Wait(TimeSpan.FromSeconds(10))) + var events = await WaitForEvents(_projectionNamesBuilder.GetEmittedStreamsName(), _numberOfTrackedEvents); + if (events.Length != _numberOfTrackedEvents) { Assert.Fail("Timed out waiting for emitted streams"); } - - var emittedStreamResult = - await _conn.ReadStreamEventsForwardAsync(_projectionNamesBuilder.GetEmittedStreamsName(), 0, - _numberOfTrackedEvents, false, _credentials); - Assert.AreEqual(_numberOfTrackedEvents, emittedStreamResult.Events.Length); - Assert.AreEqual(SliceReadStatus.Success, emittedStreamResult.Status); + while (_eventAppeared.CurrentCount > 0) + _eventAppeared.Signal(); } protected override Task When() @@ -74,9 +61,8 @@ public async Task should_have_deleted_the_tracked_emitted_streams() { for (int i = 0; i < _numberOfTrackedEvents; i++) { - var result = await _conn.ReadStreamEventsForwardAsync(String.Format(_testStreamFormat, i), 0, 1, false, - new EventStore.ClientAPI.SystemData.UserCredentials("admin", "changeit")); - Assert.AreEqual(SliceReadStatus.StreamNotFound, result.Status); + var result = await ReadEvents(String.Format(_testStreamFormat, i), 1); + Assert.That(result.Exists, Is.False); } } @@ -84,16 +70,14 @@ public async Task should_have_deleted_the_tracked_emitted_streams() [Test] public async Task should_have_deleted_the_checkpoint_stream() { - var result = await _conn.ReadStreamEventsForwardAsync(_projectionNamesBuilder.GetEmittedStreamsCheckpointName(), - 0, 1, false, new EventStore.ClientAPI.SystemData.UserCredentials("admin", "changeit")); - Assert.AreEqual(SliceReadStatus.StreamNotFound, result.Status); + var result = await ReadEvents(_projectionNamesBuilder.GetEmittedStreamsCheckpointName(), 1); + Assert.That(result.Exists, Is.False); } [Test] public async Task should_have_deleted_the_emitted_streams_stream() { - var result = await _conn.ReadStreamEventsForwardAsync(_projectionNamesBuilder.GetEmittedStreamsName(), 0, 1, - false, new EventStore.ClientAPI.SystemData.UserCredentials("admin", "changeit")); - Assert.AreEqual(SliceReadStatus.StreamNotFound, result.Status); + var result = await ReadEvents(_projectionNamesBuilder.GetEmittedStreamsName(), 1); + Assert.That(result.Exists, Is.False); } } diff --git a/src/EventStore.Projections.Core.Tests/Services/emitted_streams_tracker/when_tracking/with_tracking_disabled.cs b/src/EventStore.Projections.Core.Tests/Services/emitted_streams_tracker/when_tracking/with_tracking_disabled.cs index 9bc197b48e..892dfd8083 100644 --- a/src/EventStore.Projections.Core.Tests/Services/emitted_streams_tracker/when_tracking/with_tracking_disabled.cs +++ b/src/EventStore.Projections.Core.Tests/Services/emitted_streams_tracker/when_tracking/with_tracking_disabled.cs @@ -1,12 +1,12 @@ using System; using System.Threading; using System.Threading.Tasks; -using EventStore.ClientAPI.SystemData; using EventStore.Core.Tests; using EventStore.Projections.Core.Services.Processing; using EventStore.Projections.Core.Services.Processing.Checkpointing; using EventStore.Projections.Core.Services.Processing.Emitting; using EventStore.Projections.Core.Services.Processing.Emitting.EmittedEvents; +using Grpc.Core; using NUnit.Framework; namespace EventStore.Projections.Core.Tests.Services.emitted_streams_tracker.when_tracking; @@ -15,7 +15,6 @@ namespace EventStore.Projections.Core.Tests.Services.emitted_streams_tracker.whe public class with_tracking_disabled : SpecificationWithEmittedStreamsTrackerAndDeleter { private CountdownEvent _eventAppeared = new CountdownEvent(1); - private UserCredentials _credentials = new UserCredentials("admin", "changeit"); protected override TimeSpan Timeout { get; } = TimeSpan.FromSeconds(10); @@ -27,11 +26,15 @@ protected override Task Given() protected override async Task When() { - var sub = await _conn.SubscribeToStreamAsync(_projectionNamesBuilder.GetEmittedStreamsName(), true, (s, evnt) => - { - _eventAppeared.Signal(); - return Task.CompletedTask; - }, userCredentials: _credentials); + using var observation = new CancellationTokenSource(TimeSpan.FromSeconds(5)); + using var subscription = SubscribeToStream( + _projectionNamesBuilder.GetEmittedStreamsName(), + true, + observation.Token); + Assert.That(await subscription.ResponseStream.MoveNext(observation.Token), Is.True); + Assert.That( + subscription.ResponseStream.Current.ContentCase, + Is.EqualTo(EventStore.Client.Streams.ReadResp.ContentOneofCase.Confirmation)); _emittedStreamsTracker.TrackEmittedStream(new EmittedEvent[] { new EmittedDataEvent( @@ -39,16 +42,32 @@ protected override async Task When() "data", null, CheckpointTag.FromPosition(0, 100, 50), null, null) }); - _eventAppeared.Wait(TimeSpan.FromSeconds(5)); - sub.Unsubscribe(); + try + { + while (await subscription.ResponseStream.MoveNext(observation.Token)) + { + if (subscription.ResponseStream.Current.ContentCase == + EventStore.Client.Streams.ReadResp.ContentOneofCase.Event) + { + _eventAppeared.Signal(); + break; + } + } + } + catch (OperationCanceledException) when (observation.IsCancellationRequested) + { + } + catch (RpcException ex) when (observation.IsCancellationRequested && + ex.StatusCode is StatusCode.Cancelled or StatusCode.DeadlineExceeded) + { + } } [Test] public async Task should_write_a_stream_tracked_event() { - var result = await _conn.ReadStreamEventsForwardAsync(_projectionNamesBuilder.GetEmittedStreamsName(), 0, 200, - false, _credentials); + var result = await ReadEvents(_projectionNamesBuilder.GetEmittedStreamsName(), 200); Assert.AreEqual(0, result.Events.Length); - Assert.AreEqual(1, _eventAppeared.CurrentCount); //no event appeared should get through + Assert.AreEqual(1, _eventAppeared.CurrentCount); } } diff --git a/src/EventStore.Projections.Core.Tests/Services/emitted_streams_tracker/when_tracking/with_tracking_enabled.cs b/src/EventStore.Projections.Core.Tests/Services/emitted_streams_tracker/when_tracking/with_tracking_enabled.cs index aaabcce695..8634ecc5ae 100644 --- a/src/EventStore.Projections.Core.Tests/Services/emitted_streams_tracker/when_tracking/with_tracking_enabled.cs +++ b/src/EventStore.Projections.Core.Tests/Services/emitted_streams_tracker/when_tracking/with_tracking_enabled.cs @@ -1,8 +1,7 @@ using System; using System.Threading; using System.Threading.Tasks; -using EventStore.ClientAPI.Common.Utils; -using EventStore.ClientAPI.SystemData; +using EventStore.Common.Utils; using EventStore.Core.Tests; using EventStore.Projections.Core.Services.Processing; using EventStore.Projections.Core.Services.Processing.Checkpointing; @@ -16,33 +15,26 @@ namespace EventStore.Projections.Core.Tests.Services.emitted_streams_tracker.whe public class with_tracking_enabled : SpecificationWithEmittedStreamsTrackerAndDeleter { private CountdownEvent _eventAppeared = new CountdownEvent(1); - private UserCredentials _credentials = new UserCredentials("admin", "changeit"); protected override async Task When() { - var sub = await _conn.SubscribeToStreamAsync(_projectionNamesBuilder.GetEmittedStreamsName(), true, (s, evnt) => - { - _eventAppeared.Signal(); - return Task.CompletedTask; - }, userCredentials: _credentials); - _emittedStreamsTracker.TrackEmittedStream(new EmittedEvent[] { new EmittedDataEvent( "test_stream", Guid.NewGuid(), "type1", true, "data", null, CheckpointTag.FromPosition(0, 100, 50), null, null) }); - _eventAppeared.Wait(TimeSpan.FromSeconds(5)); - sub.Unsubscribe(); + var events = await WaitForEvents(_projectionNamesBuilder.GetEmittedStreamsName(), 1); + if (events.Length == 1) + _eventAppeared.Signal(); } [Test] public async Task should_write_a_stream_tracked_event() { - var result = await _conn.ReadStreamEventsForwardAsync(_projectionNamesBuilder.GetEmittedStreamsName(), 0, 200, - false, _credentials); + var result = await ReadEvents(_projectionNamesBuilder.GetEmittedStreamsName(), 200); Assert.AreEqual(1, result.Events.Length); - Assert.AreEqual("test_stream", Helper.UTF8NoBom.GetString(result.Events[0].Event.Data)); + Assert.AreEqual("test_stream", Helper.UTF8NoBom.GetString(result.Events[0].Event.Data.ToByteArray())); Assert.AreEqual(0, _eventAppeared.CurrentCount); } } diff --git a/src/EventStore.Projections.Core.Tests/Services/emitted_streams_tracker/when_tracking/with_tracking_enabled_with_duplicate_event_streams.cs b/src/EventStore.Projections.Core.Tests/Services/emitted_streams_tracker/when_tracking/with_tracking_enabled_with_duplicate_event_streams.cs index f83e80b5ca..df6fca762b 100644 --- a/src/EventStore.Projections.Core.Tests/Services/emitted_streams_tracker/when_tracking/with_tracking_enabled_with_duplicate_event_streams.cs +++ b/src/EventStore.Projections.Core.Tests/Services/emitted_streams_tracker/when_tracking/with_tracking_enabled_with_duplicate_event_streams.cs @@ -1,8 +1,7 @@ using System; using System.Threading; using System.Threading.Tasks; -using EventStore.ClientAPI.Common.Utils; -using EventStore.ClientAPI.SystemData; +using EventStore.Common.Utils; using EventStore.Core.Tests; using EventStore.Projections.Core.Services.Processing; using EventStore.Projections.Core.Services.Processing.Checkpointing; @@ -16,18 +15,11 @@ namespace EventStore.Projections.Core.Tests.Services.emitted_stream_manager.when public class with_tracking_enabled_with_duplicate_event_streams : SpecificationWithEmittedStreamsTrackerAndDeleter { private CountdownEvent _eventAppeared = new CountdownEvent(2); - private UserCredentials _credentials = new UserCredentials("admin", "changeit"); protected override TimeSpan Timeout { get; } = TimeSpan.FromSeconds(10); protected override async Task When() { - var sub = await _conn.SubscribeToStreamAsync(_projectionNamesBuilder.GetEmittedStreamsName(), true, (s, evnt) => - { - _eventAppeared.Signal(); - return Task.CompletedTask; - }, userCredentials: _credentials); - _emittedStreamsTracker.TrackEmittedStream(new EmittedEvent[] { new EmittedDataEvent( "test_stream", Guid.NewGuid(), "type1", true, @@ -37,17 +29,17 @@ protected override async Task When() "data", null, CheckpointTag.FromPosition(0, 100, 50), null, null) }); - _eventAppeared.Wait(TimeSpan.FromSeconds(5)); - sub.Unsubscribe(); + var events = await WaitForEvents(_projectionNamesBuilder.GetEmittedStreamsName(), 1); + if (events.Length == 1) + _eventAppeared.Signal(); } [Test] public async Task should_at_best_attempt_to_track_a_unique_list_of_streams() { - var result = await _conn.ReadStreamEventsForwardAsync(_projectionNamesBuilder.GetEmittedStreamsName(), 0, 200, - false, _credentials); + var result = await ReadEvents(_projectionNamesBuilder.GetEmittedStreamsName(), 200); Assert.AreEqual(1, result.Events.Length); - Assert.AreEqual("test_stream", Helper.UTF8NoBom.GetString(result.Events[0].Event.Data)); + Assert.AreEqual("test_stream", Helper.UTF8NoBom.GetString(result.Events[0].Event.Data.ToByteArray())); Assert.AreEqual(1, _eventAppeared.CurrentCount); //only 1 event appeared should get through } } diff --git a/src/EventStore.Projections.Core.Tests/Services/event_filter/include_everything_event_filter.cs b/src/EventStore.Projections.Core.Tests/Services/event_filter/include_everything_event_filter.cs index d463501bcd..9cb6586c68 100644 --- a/src/EventStore.Projections.Core.Tests/Services/event_filter/include_everything_event_filter.cs +++ b/src/EventStore.Projections.Core.Tests/Services/event_filter/include_everything_event_filter.cs @@ -1,4 +1,4 @@ -using EventStore.ClientAPI.Common; +using EventStore.Core.Services; using NUnit.Framework; namespace EventStore.Projections.Core.Tests.Services.event_filter; diff --git a/src/EventStore.Projections.Core.Tests/Services/event_filter/include_everything_handling_deleted_notifications_event_filter.cs b/src/EventStore.Projections.Core.Tests/Services/event_filter/include_everything_handling_deleted_notifications_event_filter.cs index 4f387d57be..255ef701ee 100644 --- a/src/EventStore.Projections.Core.Tests/Services/event_filter/include_everything_handling_deleted_notifications_event_filter.cs +++ b/src/EventStore.Projections.Core.Tests/Services/event_filter/include_everything_handling_deleted_notifications_event_filter.cs @@ -1,4 +1,4 @@ -using EventStore.ClientAPI.Common; +using EventStore.Core.Services; using NUnit.Framework; namespace EventStore.Projections.Core.Tests.Services.event_filter; diff --git a/src/EventStore.Projections.Core.Tests/Services/grpc_service/SpecificationWithNodeAndProjectionSubsystem.cs b/src/EventStore.Projections.Core.Tests/Services/grpc_service/SpecificationWithNodeAndProjectionSubsystem.cs index c607a0b792..be7d588a72 100644 --- a/src/EventStore.Projections.Core.Tests/Services/grpc_service/SpecificationWithNodeAndProjectionSubsystem.cs +++ b/src/EventStore.Projections.Core.Tests/Services/grpc_service/SpecificationWithNodeAndProjectionSubsystem.cs @@ -1,24 +1,26 @@ using System; using System.Text; using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.ClientAPI.SystemData; +using EventStore.Client.Streams; using EventStore.Common.Options; -using EventStore.Core.Services; +using EventStore.Core.Services.Transport.Grpc; using EventStore.Core.Tests; -using EventStore.Core.Tests.ClientAPI.Helpers; using EventStore.Core.Tests.Helpers; using EventStore.Core.Util; using EventStore.Projections.Core.Services.Processing; +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.Projections.Core.Tests.Services.grpc_service; public abstract class SpecificationWithNodeAndProjectionSubsystem : SpecificationWithDirectoryPerTestFixture { protected MiniNode _node; - protected IEventStoreConnection _connection; - protected UserCredentials _credentials; + private GrpcChannel _channel; + protected StreamsClient _connection; protected TimeSpan _timeout; protected string _tag; protected virtual TimeSpan StartupTimeout => TimeSpan.FromMinutes(5); @@ -30,20 +32,18 @@ public abstract class SpecificationWithNodeAndProjectionSubsystem TestConnection.CreateMiniNodeClient(_node.TcpEndPoint), - connection => connection.ReadAllEventsForwardAsync(Position.Start, 1, false, _credentials), - StartupTimeout); + _channel = GrpcChannel.ForAddress(new UriBuilder { Scheme = Uri.UriSchemeHttps }.Uri, + new GrpcChannelOptions { HttpClient = _node.HttpClient, DisposeHttpClient = false }); + _connection = new StreamsClient(_channel); try { @@ -67,22 +67,7 @@ public override async Task TestFixtureSetUp() [OneTimeTearDown] public override async Task TestFixtureTearDown() { - if (_connection != null) - { - try - { - await TestConnectionLifecycle.CloseConnectionAndWait(_connection, _timeout); - } - catch - { - TestConnectionLifecycle.TryCloseConnection(_connection); - } - finally - { - TestConnectionLifecycle.DisposeIfNeeded(_connection); - } - } - + _channel?.Dispose(); await _node.Shutdown(); await Task.Delay(1000); @@ -101,14 +86,32 @@ protected MiniNode CreateNode() subsystems: [_projectionsSubsystem]); } - protected EventData CreateEvent(string eventType, string data) + protected async Task PostEvent(string stream, string eventType, string data) { - return new EventData(Guid.NewGuid(), eventType, true, Encoding.UTF8.GetBytes(data), null); - } - - protected Task PostEvent(string stream, string eventType, string data) - { - return _connection.AppendToStreamAsync(stream, ExpectedVersion.Any, new[] { CreateEvent(eventType, data) }); + using var call = _connection.Append(); + await call.RequestStream.WriteAsync(new AppendReq + { + Options = new() + { + Any = new(), + StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(stream) } + } + }); + await call.RequestStream.WriteAsync(new AppendReq + { + ProposedMessage = new() + { + Id = Uuid.NewUuid().ToDto(), + Data = ByteString.CopyFromUtf8(data), + CustomMetadata = ByteString.Empty, + Metadata = { + { GrpcMetadata.Type, eventType }, + { GrpcMetadata.ContentType, GrpcMetadata.ContentTypes.ApplicationJson } + } + } + }); + await call.RequestStream.CompleteAsync(); + await call.ResponseAsync; } protected string CreateStandardQuery(string stream) diff --git a/src/EventStore.Projections.Core.Tests/Services/projections_manager/when_deleting_a_system_projection.cs b/src/EventStore.Projections.Core.Tests/Services/projections_manager/when_deleting_a_system_projection.cs index 2b46f51a4c..b4c19cba97 100644 --- a/src/EventStore.Projections.Core.Tests/Services/projections_manager/when_deleting_a_system_projection.cs +++ b/src/EventStore.Projections.Core.Tests/Services/projections_manager/when_deleting_a_system_projection.cs @@ -2,7 +2,7 @@ using System.Collections; using System.Collections.Generic; using System.Linq; -using EventStore.ClientAPI.Common.Utils; +using EventStore.Common.Utils; using EventStore.Core.Messages; using EventStore.Core.Messaging; using EventStore.Core.Tests; diff --git a/src/EventStore.Projections.Core/Services/Management/ManagedProjection.cs b/src/EventStore.Projections.Core/Services/Management/ManagedProjection.cs index 40a8abb7c8..ab131a91eb 100644 --- a/src/EventStore.Projections.Core/Services/Management/ManagedProjection.cs +++ b/src/EventStore.Projections.Core/Services/Management/ManagedProjection.cs @@ -1015,7 +1015,7 @@ private void DeleteStreamCompleted(ClientMessage.DeleteStreamCompleted message, Action completed) { // currently, WrongExpectedVersion is returned when deleting non-existing streams, even when specifying ExpectedVersion.Any. - // it is not too intuitive but changing the response would break the contract and compatibility with TCP/gRPC/web clients or require adding a new error code to all clients. + // Changing this response requires a coordinated public API and Admin UI contract change. // note: we don't need to check if CurrentVersion == -1 here to make sure it's a non-existing stream since the deletion is done with ExpectedVersion.Any if (message.Result == OperationResult.WrongExpectedVersion) { diff --git a/src/EventStore.Projections.Core/Services/Processing/Emitting/EmittedStreamsDeleter.cs b/src/EventStore.Projections.Core/Services/Processing/Emitting/EmittedStreamsDeleter.cs index 0ecd944c9c..24b2a9460f 100644 --- a/src/EventStore.Projections.Core/Services/Processing/Emitting/EmittedStreamsDeleter.cs +++ b/src/EventStore.Projections.Core/Services/Processing/Emitting/EmittedStreamsDeleter.cs @@ -86,7 +86,7 @@ private void ReadCompleted(ClientMessage.ReadStreamEventsForwardCompleted onRead SystemAccounts.System, x => { // currently, WrongExpectedVersion is returned when deleting non-existing streams, even when specifying ExpectedVersion.Any. - // it is not too intuitive but changing the response would break the contract and compatibility with TCP/gRPC/web clients or require adding a new error code to all clients. + // Changing this response requires a coordinated public API and Admin UI contract change. // note: we don't need to check if CurrentVersion == -1 here to make sure it's a non-existing stream since the deletion is done with ExpectedVersion.Any if (x.Result == OperationResult.WrongExpectedVersion) { @@ -108,7 +108,7 @@ private void ReadCompleted(ClientMessage.ReadStreamEventsForwardCompleted onRead SystemAccounts.System, y => { // currently, WrongExpectedVersion is returned when deleting non-existing streams, even when specifying ExpectedVersion.Any. - // it is not too intuitive but changing the response would break the contract and compatibility with TCP/gRPC/web clients or require adding a new error code to all clients. + // Changing this response requires a coordinated public API and Admin UI contract change. // note: we don't need to check if CurrentVersion == -1 here to make sure it's a non-existing stream since the deletion is done with ExpectedVersion.Any if (x.Result == OperationResult.WrongExpectedVersion) {