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/Services/SpecificationWithEmittedStreamsTrackerAndDeleter.cs b/src/EventStore.Projections.Core.Tests/Services/SpecificationWithEmittedStreamsTrackerAndDeleter.cs
index 22ae22be66..a804f14d7b 100644
--- a/src/EventStore.Projections.Core.Tests/Services/SpecificationWithEmittedStreamsTrackerAndDeleter.cs
+++ b/src/EventStore.Projections.Core.Tests/Services/SpecificationWithEmittedStreamsTrackerAndDeleter.cs
@@ -1,15 +1,33 @@
+using System;
+using System.Collections.Generic;
+using System.IO;
+using System.Linq;
+using System.Text;
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
{
+ private GrpcChannel _channel;
+ protected MiniNode _node;
+ protected StreamsClient _client;
protected IEmittedStreamsTracker _emittedStreamsTracker;
protected IEmittedStreamsDeleter _emittedStreamsDeleter;
protected ProjectionNamesBuilder _projectionNamesBuilder;
@@ -17,8 +35,32 @@ public abstract class SpecificationWithEmittedStreamsTrackerAndDeleter(PathName);
+ await _node.Start();
+ _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 +80,79 @@ 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 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());
+ var events = new List();
+ while (await call.ResponseStream.MoveNext(default))
+ if (call.ResponseStream.Current.Event is { } resolvedEvent)
+ events.Add(resolvedEvent);
+ return events.ToArray();
+ }
+
+ protected async Task WaitForEvents(string stream, int count)
+ {
+ var deadline = DateTime.UtcNow + Timeout;
+ ReadEvent[] events;
+ do
+ {
+ events = await ReadEvents(stream, count);
+ if (events.Length >= count)
+ return events;
+ await Task.Delay(50);
+ } while (DateTime.UtcNow < deadline);
+
+ return events;
+ }
+
+ private static CallOptions AdminCallOptions() => new(
+ credentials: CallCredentials.FromInterceptor((_, metadata) =>
+ {
+ metadata.Add("authorization",
+ $"Basic {Convert.ToBase64String(Encoding.ASCII.GetBytes("admin:changeit"))}");
+ return Task.CompletedTask;
+ }));
}
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..8802e0db6c 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 events = await ReadEvents(_testStreamName, 1);
+ Assert.AreEqual(0, events.Length);
}
[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 events = await ReadEvents(_projectionNamesBuilder.GetEmittedStreamsCheckpointName(), 1);
+ Assert.AreEqual(0, events.Length);
}
[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 events = await ReadEvents(_projectionNamesBuilder.GetEmittedStreamsName(), 1);
+ Assert.AreEqual(0, events.Length);
}
}
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..e984a4d928 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 events = await ReadEvents(String.Format(_testStreamFormat, i), 1);
+ Assert.AreEqual(0, events.Length);
}
}
@@ -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 events = await ReadEvents(_projectionNamesBuilder.GetEmittedStreamsCheckpointName(), 1);
+ Assert.AreEqual(0, events.Length);
}
[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 events = await ReadEvents(_projectionNamesBuilder.GetEmittedStreamsName(), 1);
+ Assert.AreEqual(0, events.Length);
}
}
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..4bc79bcc51 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,7 +1,6 @@
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;
@@ -15,7 +14,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,28 +25,20 @@ 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);
-
_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();
+ await Task.Delay(100);
}
[Test]
public async Task should_write_a_stream_tracked_event()
{
- var result = await _conn.ReadStreamEventsForwardAsync(_projectionNamesBuilder.GetEmittedStreamsName(), 0, 200,
- false, _credentials);
- Assert.AreEqual(0, result.Events.Length);
+ var events = await ReadEvents(_projectionNamesBuilder.GetEmittedStreamsName(), 200);
+ Assert.AreEqual(0, events.Length);
Assert.AreEqual(1, _eventAppeared.CurrentCount); //no event appeared should get through
}
}
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..97ef20b960 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);
- Assert.AreEqual(1, result.Events.Length);
- Assert.AreEqual("test_stream", Helper.UTF8NoBom.GetString(result.Events[0].Event.Data));
+ var events = await ReadEvents(_projectionNamesBuilder.GetEmittedStreamsName(), 200);
+ Assert.AreEqual(1, events.Length);
+ Assert.AreEqual("test_stream", Helper.UTF8NoBom.GetString(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..49e5817a99 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);
- Assert.AreEqual(1, result.Events.Length);
- Assert.AreEqual("test_stream", Helper.UTF8NoBom.GetString(result.Events[0].Event.Data));
+ var events = await ReadEvents(_projectionNamesBuilder.GetEmittedStreamsName(), 200);
+ Assert.AreEqual(1, events.Length);
+ Assert.AreEqual("test_stream", Helper.UTF8NoBom.GetString(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..93c6e33ffe 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,17 @@ 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 +66,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 +85,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)
{