diff --git a/.github/workflows/build-container-ubuntu-lts.yml b/.github/workflows/build-container-ubuntu-lts.yml index 91ffc103fb..bbee6ddc43 100644 --- a/.github/workflows/build-container-ubuntu-lts.yml +++ b/.github/workflows/build-container-ubuntu-lts.yml @@ -35,7 +35,6 @@ jobs: fail-fast: true matrix: include: - - test-group-name: core-clientapi-persistent - test-group-name: core-clientapi-security - test-group-name: core-clientapi-streams - test-group-name: core-http diff --git a/scripts/test.sh b/scripts/test.sh index 7e90473284..b606d3c713 100755 --- a/scripts/test.sh +++ b/scripts/test.sh @@ -65,9 +65,6 @@ load_requested_projects() { core-clientapi) requested_projects=("${core_clientapi_projects[@]}") ;; - core-clientapi-persistent) - requested_projects=("${core_clientapi_projects[@]}") - ;; core-clientapi-security) requested_projects=("${core_clientapi_projects[@]}") ;; @@ -207,9 +204,6 @@ project_filter() { core-clientapi:EventStore.Core.Tests) printf '%s\n' "FullyQualifiedName~EventStore.Core.Tests.ClientAPI" ;; - core-clientapi-persistent:EventStore.Core.Tests) - printf '%s\n' "FullyQualifiedName~EventStore.Core.Tests.ClientAPI&(FullyQualifiedName~persistent|FullyQualifiedName~Persistent)&FullyQualifiedName!~EventStore.Core.Tests.ClientAPI.Security" - ;; core-clientapi-security:EventStore.Core.Tests) printf '%s\n' "FullyQualifiedName~EventStore.Core.Tests.ClientAPI.Security" ;; @@ -247,9 +241,6 @@ project_timeout() { core-clientapi:EventStore.Core.Tests) printf '%s\n' "20m" ;; - core-clientapi-persistent:EventStore.Core.Tests) - printf '%s\n' "20m" - ;; core-clientapi-security:EventStore.Core.Tests) printf '%s\n' "20m" ;; diff --git a/src/EventStore.Core.Tests/ClientAPI/ExpectedVersion64Bit/persistent_subscription_with_event_numbers_greater_than_2_billion.cs b/src/EventStore.Core.Tests/ClientAPI/ExpectedVersion64Bit/persistent_subscription_with_event_numbers_greater_than_2_billion.cs deleted file mode 100644 index 6eeacf86b2..0000000000 --- a/src/EventStore.Core.Tests/ClientAPI/ExpectedVersion64Bit/persistent_subscription_with_event_numbers_greater_than_2_billion.cs +++ /dev/null @@ -1,80 +0,0 @@ -using System; -using System.Collections.Generic; -using System.Threading; -using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.Core.Data; -using NUnit.Framework; - -namespace EventStore.Core.Tests.ClientAPI.ExpectedVersion64Bit; - -[TestFixture(typeof(LogFormat.V2), typeof(string))] -[Category("ClientAPI"), Category("LongRunning")] -public class PersistentSubscriptionWithEventNumbersGreaterThan2Billion - : MiniNodeWithExistingRecords -{ - private const long intMaxValue = (long)int.MaxValue; - - private string _streamId = "persistent-subscription-stream"; - - private EventRecord _r1, _r2; - - public override async ValueTask WriteTestScenario(CancellationToken token) - { - _r1 = await WriteSingleEvent(_streamId, intMaxValue + 1, new string('.', 3000), token: token); - _r2 = await WriteSingleEvent(_streamId, intMaxValue + 2, new string('.', 3000), token: token); - } - - public override async Task Given() - { - _store = BuildConnection(Node); - await _store.ConnectAsync(); - await _store.SetStreamMetadataAsync(_streamId, EventStore.ClientAPI.ExpectedVersion.Any, - EventStore.ClientAPI.StreamMetadata.Create(truncateBefore: intMaxValue + 1)); - } - - [Test] - public async Task should_be_able_to_create_the_persistent_subscription() - { - var groupId = "group-" + Guid.NewGuid().ToString(); - var settings = PersistentSubscriptionSettings.Create().StartFrom(intMaxValue); - await _store.CreatePersistentSubscriptionAsync(_streamId, groupId, settings, DefaultData.AdminCredentials); - } - - [Test] - public async Task should_be_able_to_update_the_persistent_subscription() - { - var groupId = "group-" + Guid.NewGuid().ToString(); - var settings = PersistentSubscriptionSettings.Create(); - await _store.CreatePersistentSubscriptionAsync(_streamId, groupId, settings, DefaultData.AdminCredentials); - - settings = PersistentSubscriptionSettings.Create().StartFrom(intMaxValue); - await _store.UpdatePersistentSubscriptionAsync(_streamId, groupId, settings, DefaultData.AdminCredentials); - } - - [Test] - public async Task should_be_able_to_connect_to_persistent_subscription() - { - var groupId = "group-" + Guid.NewGuid().ToString(); - var settings = PersistentSubscriptionSettings.Create().StartFrom(intMaxValue); - await _store.CreatePersistentSubscriptionAsync(_streamId, groupId, settings, DefaultData.AdminCredentials); - - var evnt = new EventData(Guid.NewGuid(), "EventType", false, new byte[10], new byte[15]); - List receivedEvents = new List(); - var countdown = new CountdownEvent(3); - await _store.ConnectToPersistentSubscriptionAsync(_streamId, groupId, (s, e) => - { - receivedEvents.Add(e); - countdown.Signal(); - return Task.CompletedTask; - }, userCredentials: DefaultData.AdminCredentials); - - await _store.AppendToStreamAsync(_streamId, intMaxValue + 2, evnt); - - Assert.That(countdown.Wait(TimeSpan.FromSeconds(5)), "Timed out waiting for events to appear"); - - Assert.AreEqual(_r1.EventId, receivedEvents[0].Event.EventId); - Assert.AreEqual(_r2.EventId, receivedEvents[1].Event.EventId); - Assert.AreEqual(evnt.EventId, receivedEvents[2].Event.EventId); - } -} diff --git a/src/EventStore.Core.Tests/ClientAPI/connecting_to_a_persistent_subscription.cs b/src/EventStore.Core.Tests/ClientAPI/connecting_to_a_persistent_subscription.cs deleted file mode 100644 index a5edc279b4..0000000000 --- a/src/EventStore.Core.Tests/ClientAPI/connecting_to_a_persistent_subscription.cs +++ /dev/null @@ -1,980 +0,0 @@ -using System; -using System.Collections.Generic; -using System.Text; -using System.Threading; -using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.ClientAPI.ClientOperations; -using EventStore.ClientAPI.Common; -using EventStore.ClientAPI.Common.Utils; -using EventStore.ClientAPI.Exceptions; -using NUnit.Framework; - -namespace EventStore.Core.Tests.ClientAPI; - -[Category("LongRunning"), Category("ClientAPI")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class connect_to_non_existing_persistent_subscription_with_permissions - : SpecificationWithMiniNode -{ - private Exception _caught; - - protected override Task When() - { - _caught = Assert.Throws( - () => - { - _conn.ConnectToPersistentSubscription( - "nonexisting2", - "foo", - (sub, e) => - { - Console.Write("appeared"); - return Task.CompletedTask; - }, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials); - throw new Exception("should have thrown"); - }).InnerException; - return Task.CompletedTask; - } - - [Test] - public void the_completion_fails() - { - Assert.IsNotNull(_caught); - } - - [Test] - public void the_exception_is_an_argument_exception() - { - Assert.IsInstanceOf(_caught); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class connect_to_existing_persistent_subscription_with_permissions : SpecificationWithMiniNode -{ - private EventStorePersistentSubscriptionBase _sub; - private readonly string _stream = Guid.NewGuid().ToString(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - protected override async Task Given() - { - await _conn.CreatePersistentSubscriptionAsync(_stream, "agroupname17", _settings, - DefaultData.AdminCredentials); - } - - protected override Task When() - { - _sub = _conn.ConnectToPersistentSubscription(_stream, - "agroupname17", - (sub, e) => - { - Console.Write("appeared"); - return Task.CompletedTask; - }, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials); - - return Task.CompletedTask; - } - - [Test] - public void the_subscription_succeeds() - { - Assert.IsNotNull(_sub); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class connect_to_existing_persistent_subscription_without_permissions : SpecificationWithMiniNode -{ - private readonly string _stream = "$" + Guid.NewGuid(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - protected override Task When() - { - return _conn.CreatePersistentSubscriptionAsync(_stream, "agroupname55", _settings, - DefaultData.AdminCredentials); - } - - [Test] - public void the_subscription_fails_to_connect() - { - try - { - _conn.ConnectToPersistentSubscription( - _stream, - "agroupname55", - (sub, e) => - { - Console.Write("appeared"); - return Task.CompletedTask; - }, - (sub, reason, ex) => Console.WriteLine("dropped.")); - throw new Exception("should have thrown."); - } - catch (Exception ex) - { - Assert.IsInstanceOf(ex); - Assert.IsInstanceOf(ex.InnerException); - } - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class connect_to_existing_persistent_subscription_with_max_one_client : SpecificationWithMiniNode -{ - private readonly string _stream = "$" + Guid.NewGuid(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent() - .WithMaxSubscriberCountOf(1); - - private Exception _exception; - - private const string _group = "startinbeginning1"; - - protected override async Task Given() - { - await base.Given(); - await _conn.CreatePersistentSubscriptionAsync(_stream, _group, _settings, - DefaultData.AdminCredentials); - _conn.ConnectToPersistentSubscription( - _stream, - _group, - (s, e) => - { - s.Acknowledge(e); - return Task.CompletedTask; - }, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials); - } - - protected override Task When() - { - _exception = Assert.Throws(() => - { - _conn.ConnectToPersistentSubscription( - _stream, - _group, - (s, e) => - { - s.Acknowledge(e); - return Task.CompletedTask; - }, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials); - throw new Exception("should have thrown."); - }).InnerException; - return Task.CompletedTask; - } - - [Test] - public void the_second_subscription_fails_to_connect() - { - Assert.IsInstanceOf(_exception); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class - connect_to_existing_persistent_subscription_with_start_from_beginning_and_no_stream : - SpecificationWithMiniNode -{ - private readonly string _stream = "$" + Guid.NewGuid(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromBeginning(); - - private readonly Guid _id = Guid.NewGuid(); - private readonly TaskCompletionSource _firstEventSource = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); - - private const string _group = "startinbeginning1"; - - protected override async Task Given() - { - await _conn.CreatePersistentSubscriptionAsync(_stream, _group, _settings, - DefaultData.AdminCredentials); - _conn.ConnectToPersistentSubscription( - _stream, - _group, - HandleEvent, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials); - } - - protected override Task When() - { - return _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(_id, "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), new byte[0])); - } - - private Task HandleEvent(EventStorePersistentSubscriptionBase sub, ResolvedEvent resolvedEvent) - { - _firstEventSource.TrySetResult(resolvedEvent); - return Task.CompletedTask; - } - - [Test] - public async Task the_subscription_gets_event_zero_as_its_first_event() - { - var firstEvent = await _firstEventSource.Task.WithTimeout(TimeSpan.FromSeconds(10)); - Assert.AreEqual(0, firstEvent.Event.EventNumber); - Assert.AreEqual(_id, firstEvent.Event.EventId); - } -} - - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class - connect_to_existing_persistent_subscription_with_start_from_two_and_no_stream - : SpecificationWithMiniNode -{ - private readonly string _stream = "$" + Guid.NewGuid(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFrom(2); - - private readonly AutoResetEvent _resetEvent = new AutoResetEvent(false); - private readonly Guid _id = Guid.NewGuid(); - private readonly TaskCompletionSource _firstEventSource = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); - - private const string _group = "startinbeginning1"; - - protected override async Task Given() - { - await _conn.CreatePersistentSubscriptionAsync(_stream, _group, _settings, - DefaultData.AdminCredentials); - _conn.ConnectToPersistentSubscription( - _stream, - _group, - HandleEvent, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials); - } - - protected override async Task When() - { - await _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(Guid.NewGuid(), "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), new byte[0])) -; - await _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(Guid.NewGuid(), "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), new byte[0])) -; - await _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(_id, "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), new byte[0])); - } - - private Task HandleEvent(EventStorePersistentSubscriptionBase sub, ResolvedEvent resolvedEvent) - { - _firstEventSource.TrySetResult(resolvedEvent); - return Task.CompletedTask; - } - - [Test] - public async Task the_subscription_gets_event_two_as_its_first_event() - { - var resolvedEvent = await _firstEventSource.Task.WithTimeout(TimeSpan.FromSeconds(10)); - Assert.AreEqual(2, resolvedEvent.Event.EventNumber); - Assert.AreEqual(_id, resolvedEvent.Event.EventId); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class - connect_to_existing_persistent_subscription_with_start_from_beginning_and_events_in_it : - SpecificationWithMiniNode -{ - private readonly string _stream = "$" + Guid.NewGuid(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromBeginning(); - - private readonly AutoResetEvent _resetEvent = new AutoResetEvent(false); - private ResolvedEvent _firstEvent; - private List _ids = new List(); - private bool _set = false; - - private const string _group = "startinbeginning1"; - - protected override async Task Given() - { - await WriteEvents(_conn); - await _conn.CreatePersistentSubscriptionAsync(_stream, _group, _settings, - DefaultData.AdminCredentials); - } - - private async Task WriteEvents(IEventStoreConnection connection) - { - for (int i = 0; i < 10; i++) - { - _ids.Add(Guid.NewGuid()); - await connection.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(_ids[i], "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), new byte[0])) -; - } - } - - protected override Task When() - { - _conn.ConnectToPersistentSubscription( - _stream, - _group, - HandleEvent, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials); - return Task.CompletedTask; - } - - private Task HandleEvent(EventStorePersistentSubscriptionBase sub, ResolvedEvent resolvedEvent) - { - if (!_set) - { - _set = true; - _firstEvent = resolvedEvent; - _resetEvent.Set(); - } - - return Task.CompletedTask; - } - - [Test] - public void the_subscription_gets_event_zero_as_its_first_event() - { - Assert.IsTrue(_resetEvent.WaitOne(TimeSpan.FromSeconds(10))); - Assert.AreEqual(0, _firstEvent.Event.EventNumber); - Assert.AreEqual(_ids[0], _firstEvent.Event.EventId); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class - connect_to_existing_persistent_subscription_with_start_from_beginning_not_set_and_events_in_it : - SpecificationWithMiniNode -{ - private readonly string _stream = "$" + Guid.NewGuid(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - private readonly AutoResetEvent _resetEvent = new AutoResetEvent(false); - - private const string _group = "startinbeginning1"; - - protected override async Task Given() - { - await WriteEvents(_conn); - await _conn.CreatePersistentSubscriptionAsync(_stream, _group, _settings, - DefaultData.AdminCredentials); - } - - private async Task WriteEvents(IEventStoreConnection connection) - { - for (int i = 0; i < 10; i++) - { - await connection.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(Guid.NewGuid(), "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), - new byte[0])) -; - } - } - - protected override Task When() - { - _conn.ConnectToPersistentSubscription( - _stream, - _group, - HandleEvent, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials); - return Task.CompletedTask; - } - - private Task HandleEvent(EventStorePersistentSubscriptionBase sub, ResolvedEvent resolvedEvent) - { - _resetEvent.Set(); - return Task.CompletedTask; - } - - [Test] - public void the_subscription_gets_no_events() - { - Assert.IsFalse(_resetEvent.WaitOne(TimeSpan.FromSeconds(1))); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class - connect_to_existing_persistent_subscription_with_start_from_beginning_not_set_and_events_in_it_then_event_written : - SpecificationWithMiniNode -{ - private readonly string _stream = "$" + Guid.NewGuid(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos(); - - private readonly AutoResetEvent _resetEvent = new AutoResetEvent(false); - private ResolvedEvent _firstEvent; - private Guid _id; - - private const string _group = "startinbeginning1"; - - protected override async Task Given() - { - await WriteEvents(_conn); - await _conn.CreatePersistentSubscriptionAsync(_stream, _group, _settings, - DefaultData.AdminCredentials); - _conn.ConnectToPersistentSubscription( - _stream, - _group, - HandleEvent, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials); - } - - private async Task WriteEvents(IEventStoreConnection connection) - { - for (int i = 0; i < 10; i++) - { - await connection.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(Guid.NewGuid(), "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), - new byte[0])) -; - } - } - - protected override async Task When() - { - _id = Guid.NewGuid(); - await _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(_id, "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), new byte[0])); - } - - private Task HandleEvent(EventStorePersistentSubscriptionBase sub, ResolvedEvent resolvedEvent) - { - _firstEvent = resolvedEvent; - _resetEvent.Set(); - return Task.CompletedTask; - } - - [Test] - public void the_subscription_gets_the_written_event_as_its_first_event() - { - Assert.IsTrue(_resetEvent.WaitOne(TimeSpan.FromSeconds(10))); - Assert.IsNotNull(_firstEvent); - Assert.AreEqual(10, _firstEvent.Event.EventNumber); - Assert.AreEqual(_id, _firstEvent.Event.EventId); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class - connect_to_existing_persistent_subscription_with_start_from_x_set_higher_than_x_and_events_in_it_then_event_written : - SpecificationWithMiniNode -{ - private readonly string _stream = "$" + Guid.NewGuid(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFrom(11); - - private readonly AutoResetEvent _resetEvent = new AutoResetEvent(false); - private ResolvedEvent _firstEvent; - private Guid _id; - - private const string _group = "startinbeginning1"; - - protected override async Task Given() - { - await WriteEvents(_conn); - await _conn.CreatePersistentSubscriptionAsync(_stream, _group, _settings, - DefaultData.AdminCredentials); - _conn.ConnectToPersistentSubscription( - _stream, - _group, - HandleEvent, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials); - } - - private async Task WriteEvents(IEventStoreConnection connection) - { - for (int i = 0; i < 11; i++) - { - await connection.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(Guid.NewGuid(), "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), - new byte[0])) -; - } - } - - protected override Task When() - { - _id = Guid.NewGuid(); - return _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(_id, "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), new byte[0])); - } - - private Task HandleEvent(EventStorePersistentSubscriptionBase sub, ResolvedEvent resolvedEvent) - { - _firstEvent = resolvedEvent; - _resetEvent.Set(); - return Task.CompletedTask; - } - - [Test] - public void the_subscription_gets_the_written_event_as_its_first_event() - { - Assert.IsTrue(_resetEvent.WaitOne(TimeSpan.FromSeconds(10))); - Assert.IsNotNull(_firstEvent); - Assert.AreEqual(11, _firstEvent.Event.EventNumber); - Assert.AreEqual(_id, _firstEvent.Event.EventId); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class a_nak_in_subscription_handler_in_autoack_mode_drops_the_subscription : SpecificationWithMiniNode -{ - private readonly string _stream = "$" + Guid.NewGuid(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromBeginning(); - - private readonly ManualResetEvent _resetEvent = new ManualResetEvent(false); - private Exception _exception; - private SubscriptionDropReason _reason; - - private const string _group = "naktest"; - - protected override async Task Given() - { - await _conn.CreatePersistentSubscriptionAsync(_stream, _group, _settings, - DefaultData.AdminCredentials); - _conn.ConnectToPersistentSubscription( - _stream, - _group, - HandleEvent, - Dropped, - DefaultData.AdminCredentials); - } - - private void Dropped(EventStorePersistentSubscriptionBase sub, SubscriptionDropReason reason, - Exception exception) - { - _exception = exception; - _reason = reason; - _resetEvent.Set(); - } - - protected override Task When() - { - return _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(Guid.NewGuid(), "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), new byte[0])); - } - - private static Task HandleEvent(EventStorePersistentSubscriptionBase sub, ResolvedEvent resolvedEvent) - { - throw new Exception("test"); - } - - [Test] - [Retry(5)] - public void the_subscription_gets_dropped() - { - Assert.IsTrue(_resetEvent.WaitOne(TimeSpan.FromSeconds(5))); - Assert.AreEqual(SubscriptionDropReason.EventHandlerException, _reason); - Assert.AreEqual(typeof(Exception), _exception.GetType()); - Assert.AreEqual("test", _exception.Message); - } -} - - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class - connect_to_existing_persistent_subscription_with_start_from_x_set_and_events_in_it_then_event_written : - SpecificationWithMiniNode -{ - private readonly string _stream = "$" + Guid.NewGuid(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFrom(10); - - private readonly AutoResetEvent _resetEvent = new AutoResetEvent(false); - private ResolvedEvent _firstEvent; - private Guid _id; - - private const string _group = "startinbeginning1"; - - protected override async Task Given() - { - await WriteEvents(_conn); - await _conn.CreatePersistentSubscriptionAsync(_stream, _group, _settings, - DefaultData.AdminCredentials); - _conn.ConnectToPersistentSubscription( - _stream, - _group, - HandleEvent, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials); - } - - private async Task WriteEvents(IEventStoreConnection connection) - { - for (int i = 0; i < 10; i++) - { - await connection.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(Guid.NewGuid(), "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), - new byte[0])) -; - } - } - - protected override Task When() - { - _id = Guid.NewGuid(); - return _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(_id, "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), new byte[0])); - } - - private Task HandleEvent(EventStorePersistentSubscriptionBase sub, ResolvedEvent resolvedEvent) - { - _firstEvent = resolvedEvent; - _resetEvent.Set(); - return Task.CompletedTask; - } - - [Test] - public void the_subscription_gets_the_written_event_as_its_first_event() - { - Assert.IsTrue(_resetEvent.WaitOne(TimeSpan.FromSeconds(10))); - Assert.IsNotNull(_firstEvent); - Assert.AreEqual(10, _firstEvent.Event.EventNumber); - Assert.AreEqual(_id, _firstEvent.Event.EventId); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class - connect_to_existing_persistent_subscription_with_start_from_x_set_and_events_in_it - : SpecificationWithMiniNode -{ - private readonly string _stream = "$" + Guid.NewGuid(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFrom(4); - - private readonly AutoResetEvent _resetEvent = new AutoResetEvent(false); - private ResolvedEvent _firstEvent; - private Guid _id; - - private const string _group = "startinx2"; - - protected override async Task Given() - { - await WriteEvents(_conn); - await _conn.CreatePersistentSubscriptionAsync(_stream, _group, _settings, - DefaultData.AdminCredentials); - _conn.ConnectToPersistentSubscription( - _stream, - _group, - HandleEvent, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials); - } - - private async Task WriteEvents(IEventStoreConnection connection) - { - for (int i = 0; i < 10; i++) - { - var id = Guid.NewGuid(); - await connection.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(id, "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), new byte[0])); - if (i == 4) - { - _id = id; - } - } - } - - protected override Task When() - { - return _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(_id, "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), new byte[0])); - } - - private bool _set = false; - - private Task HandleEvent(EventStorePersistentSubscriptionBase sub, ResolvedEvent resolvedEvent) - { - if (_set) - { - return Task.CompletedTask; - } - - _set = true; - _firstEvent = resolvedEvent; - _resetEvent.Set(); - return Task.CompletedTask; - } - - [Test] - public void the_subscription_gets_the_written_event_as_its_first_event() - { - Assert.IsTrue(_resetEvent.WaitOne(TimeSpan.FromSeconds(10))); - Assert.IsNotNull(_firstEvent); - Assert.AreEqual(4, _firstEvent.Event.EventNumber); - Assert.AreEqual(_id, _firstEvent.Event.EventId); - } -} - -[TestFixture(typeof(LogFormat.V2), typeof(string))] -[Category("LongRunning")] -public class - connect_to_persistent_subscription_with_link_to_event_with_event_number_greater_than_int_maxvalue : - ExpectedVersion64Bit.MiniNodeWithExistingRecords -{ - private const string StreamName = - "connect_to_persistent_subscription_with_link_to_event_with_event_number_greater_than_int_maxvalue"; - - private const long intMaxValue = (long)int.MaxValue; - - private string _linkedStreamName = "linked-" + StreamName; - private const string _group = "group-" + StreamName; - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .ResolveLinkTos() - .StartFromBeginning(); - - private readonly AutoResetEvent _resetEvent = new AutoResetEvent(false); - private ResolvedEvent _firstEvent; - private bool _set = false; - private Guid _event1Id; - - public override async ValueTask WriteTestScenario(CancellationToken token) - { - var event1 = await WriteSingleEvent(StreamName, intMaxValue + 1, new string('.', 3000), token: token); - await WriteSingleEvent(StreamName, intMaxValue + 2, new string('.', 3000), token: token); - _event1Id = event1.EventId; - } - - public override async Task Given() - { - _store = BuildConnection(Node); - await _store.ConnectAsync(); - - await _store.CreatePersistentSubscriptionAsync(_linkedStreamName, _group, _settings, - DefaultData.AdminCredentials); - _store.ConnectToPersistentSubscription( - _linkedStreamName, - _group, - HandleEvent, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials); - await _store.AppendToStreamAsync(_linkedStreamName, ExpectedVersion.Any, new EventData(Guid.NewGuid(), - SystemEventTypes.LinkTo, false, Helper.UTF8NoBom.GetBytes( - string.Format("{0}@{1}", intMaxValue + 1, StreamName)), null)); - } - - private Task HandleEvent(EventStorePersistentSubscriptionBase sub, ResolvedEvent resolvedEvent) - { - if (!_set) - { - _set = true; - _firstEvent = resolvedEvent; - _resetEvent.Set(); - } - - return Task.CompletedTask; - } - - [Test] - public void the_subscription_resolves_the_linked_event_correctly() - { - Assert.IsTrue(_resetEvent.WaitOne(TimeSpan.FromSeconds(10))); - Assert.AreEqual(intMaxValue + 1, _firstEvent.Event.EventNumber); - Assert.AreEqual(_event1Id, _firstEvent.Event.EventId); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class connect_to_persistent_subscription_with_retries : SpecificationWithMiniNode -{ - private readonly string _stream = Guid.NewGuid().ToString("N"); - protected override TimeSpan Timeout { get; } = TimeSpan.FromMinutes(5); - protected override TimeSpan StartupTimeout { get; } = TimeSpan.FromMinutes(10); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromBeginning(); - - private readonly AutoResetEvent _resetEvent = new AutoResetEvent(false); - private readonly Guid _id = Guid.NewGuid(); - int? _retryCount; - private const string _group = "retries"; - - protected override async Task Given() - { - await _conn.CreatePersistentSubscriptionAsync(_stream, _group, _settings, - DefaultData.AdminCredentials); - await _conn.ConnectToPersistentSubscriptionAsync( - _stream, - _group, - HandleEvent, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials, autoAck: false); - } - - protected override Task When() - { - return _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(_id, "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), new byte[0])); - } - - private Task HandleEvent(EventStorePersistentSubscriptionBase sub, ResolvedEvent resolvedEvent, - int? retryCount) - { - if (retryCount > 4) - { - _retryCount = retryCount; - sub.Acknowledge(resolvedEvent); - _resetEvent.Set(); - } - else - { - sub.Fail(resolvedEvent, PersistentSubscriptionNakEventAction.Retry, "Not yet tried enough times"); - } - - return Task.CompletedTask; - } - - [Test] - public void events_are_retried_until_success() - { - Assert.IsTrue(_resetEvent.WaitOne(TimeSpan.FromSeconds(10))); - Assert.AreEqual(5, _retryCount); - } -} - -//ALL - -/* - - [TestFixture, Category("LongRunning")] - public class connect_to_non_existing_persistent_all_subscription_with_permissions : SpecificationWithMiniNode - { - private Exception _caught; - - protected override void When() - { - try - { - _conn.ConnectToPersistentSubscriptionForAll("nonexisting2", - (sub, e) => Console.Write("appeared"), - (sub, reason, ex) => - { - }, - DefaultData.AdminCredentials); - throw new Exception("should have thrown"); - } - catch (Exception ex) - { - _caught = ex; - } - } - - [Test] - public void the_completion_fails() - { - Assert.IsNotNull(_caught); - } - - [Test] - public void the_exception_is_an_argument_exception() - { - Assert.IsInstanceOf(_caught.InnerException); - } - } - - [TestFixture, Category("LongRunning")] - public class connect_to_existing_persistent_all_subscription_with_permissions : SpecificationWithMiniNode - { - private EventStorePersistentSubscription _sub; - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettingsBuilder.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - protected override void When() - { - _conn.CreatePersistentSubscriptionForAllAsync("agroupname17", _settings, DefaultData.AdminCredentials).Wait(); - _sub = _conn.ConnectToPersistentSubscriptionForAll("agroupname17", - (sub, e) => Console.Write("appeared"), - (sub, reason, ex) => { }, DefaultData.AdminCredentials); - } - - [Test] - public void the_subscription_suceeds() - { - Assert.IsNotNull(_sub); - } - } - - [TestFixture, Category("LongRunning")] - public class connect_to_existing_persistent_all_subscription_without_permissions : SpecificationWithMiniNode - { - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettingsBuilder.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - protected override void When() - { - _conn.CreatePersistentSubscriptionForAllAsync("agroupname55", _settings, - DefaultData.AdminCredentials).Wait(); - } - - [Test] - public void the_subscription_fails_to_connect() - { - try - { - _conn.ConnectToPersistentSubscriptionForAll("agroupname55", - (sub, e) => Console.Write("appeared"), - (sub, reason, ex) => { }); - throw new Exception("should have thrown."); - } - catch (Exception ex) - { - Assert.IsInstanceOf(ex); - Assert.IsInstanceOf(ex.InnerException); - } - } - } -*/ diff --git a/src/EventStore.Core.Tests/ClientAPI/connecting_to_a_persistent_subscription_async.cs b/src/EventStore.Core.Tests/ClientAPI/connecting_to_a_persistent_subscription_async.cs deleted file mode 100644 index f87a979b4a..0000000000 --- a/src/EventStore.Core.Tests/ClientAPI/connecting_to_a_persistent_subscription_async.cs +++ /dev/null @@ -1,412 +0,0 @@ -using System; -using System.Collections.Generic; -using System.Text; -using System.Threading; -using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.ClientAPI.ClientOperations; -using EventStore.ClientAPI.Exceptions; -using NUnit.Framework; - -namespace EventStore.Core.Tests.ClientAPI; - -[Category("LongRunning"), Category("ClientAPI")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class connect_to_non_existing_persistent_subscription_with_permissions_async : SpecificationWithMiniNode -{ - private Exception _innerEx; - - protected override async Task When() - { - _innerEx = await AssertEx.ThrowsAsync(() => _conn.ConnectToPersistentSubscriptionAsync( - "nonexisting2", - "foo", - (sub, e) => - { - Console.Write("appeared"); - return Task.CompletedTask; - }, - (sub, reason, ex) => { }, DefaultData.AdminCredentials)); - } - - [Test] - public void the_subscription_fails_to_connect_with_argument_exception() - { - Assert.IsInstanceOf(_innerEx); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class connect_to_existing_persistent_subscription_with_permissions_async : SpecificationWithMiniNode -{ - private EventStorePersistentSubscriptionBase _sub; - private readonly string _stream = Guid.NewGuid().ToString(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - protected override async Task When() - { - await _conn.CreatePersistentSubscriptionAsync(_stream, "agroupname17", _settings, DefaultData.AdminCredentials) -; - _sub = await _conn.ConnectToPersistentSubscriptionAsync(_stream, - "agroupname17", - (sub, e) => - { - Console.Write("appeared"); - return Task.CompletedTask; - }, - (sub, reason, ex) => { }, DefaultData.AdminCredentials); - } - - [Test] - public void the_subscription_succeeds() - { - Assert.IsNotNull(_sub); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class connect_to_existing_persistent_subscription_without_permissions_async : SpecificationWithMiniNode -{ - private readonly string _stream = "$" + Guid.NewGuid(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - private Exception _innerEx; - - protected override async Task When() - { - await _conn.CreatePersistentSubscriptionAsync(_stream, "agroupname55", _settings, - DefaultData.AdminCredentials); - _innerEx = await AssertEx.ThrowsAsync(() => _conn.ConnectToPersistentSubscriptionAsync( - _stream, - "agroupname55", - (sub, e) => - { - Console.Write("appeared"); - return Task.CompletedTask; - }, - (sub, reason, ex) => Console.WriteLine("dropped."))); - } - - [Test] - public void the_subscription_fails_to_connect_with_access_denied_exception() - { - Assert.IsInstanceOf(_innerEx); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class connect_to_existing_persistent_subscription_with_max_one_client_async : SpecificationWithMiniNode -{ - private readonly string _stream = "$" + Guid.NewGuid(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent() - .WithMaxSubscriberCountOf(1); - - private Exception _innerEx; - - private const string _group = "startinbeginning1"; - private EventStorePersistentSubscriptionBase _firstConn; - - protected override async Task Given() - { - await base.Given(); - await _conn.CreatePersistentSubscriptionAsync(_stream, _group, _settings, - DefaultData.AdminCredentials); - // First connection - _firstConn = await _conn.ConnectToPersistentSubscriptionAsync( - _stream, - _group, - (s, e) => - { - s.Acknowledge(e); - return Task.CompletedTask; - }, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials); - } - - protected override async Task When() - { - _innerEx = await AssertEx.ThrowsAsync(() => - // Second connection - _conn.ConnectToPersistentSubscriptionAsync( - _stream, - _group, - (s, e) => - { - s.Acknowledge(e); - return Task.CompletedTask; - }, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials)); - } - - [Test] - public void the_first_subscription_connects_successfully() - { - Assert.IsNotNull(_firstConn); - } - - [Test] - public void the_second_subscription_throws_maximum_subscribers_reached_exception() - { - Assert.IsInstanceOf(_innerEx); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class - connect_to_existing_persistent_subscription_with_start_from_beginning_and_no_stream_async : - SpecificationWithMiniNode -{ - private readonly string _stream = "$" + Guid.NewGuid(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromBeginning(); - - private readonly AutoResetEvent _resetEvent = new AutoResetEvent(false); - private ResolvedEvent _firstEvent; - private readonly Guid _id = Guid.NewGuid(); - private bool _set = false; - - private const string _group = "startinbeginning1"; - - protected override async Task Given() - { - await _conn.CreatePersistentSubscriptionAsync(_stream, _group, _settings, - DefaultData.AdminCredentials); - - await _conn.ConnectToPersistentSubscriptionAsync( - _stream, - _group, - HandleEvent, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials); - } - - protected override Task When() - { - return _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(_id, "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), new byte[0])); - } - - private Task HandleEvent(EventStorePersistentSubscriptionBase sub, ResolvedEvent resolvedEvent) - { - if (_set) - { - return Task.CompletedTask; - } - - _set = true; - _firstEvent = resolvedEvent; - _resetEvent.Set(); - return Task.CompletedTask; - } - - [Test] - public void the_subscription_gets_event_zero_as_its_first_event() - { - Assert.IsTrue(_resetEvent.WaitOne(TimeSpan.FromSeconds(10))); - Assert.AreEqual(0, _firstEvent.Event.EventNumber); - Assert.AreEqual(_id, _firstEvent.Event.EventId); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class - connect_to_existing_persistent_subscription_with_start_from_two_and_no_stream_async : - SpecificationWithMiniNode -{ - private readonly string _stream = "$" + Guid.NewGuid(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFrom(2); - - private readonly AutoResetEvent _resetEvent = new AutoResetEvent(false); - private ResolvedEvent _firstEvent; - private readonly Guid _id = Guid.NewGuid(); - private bool _set = false; - - private const string _group = "startinbeginning1"; - - protected override async Task Given() - { - await _conn.CreatePersistentSubscriptionAsync(_stream, _group, _settings, - DefaultData.AdminCredentials); - await _conn.ConnectToPersistentSubscriptionAsync( - _stream, - _group, - HandleEvent, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials); - } - - protected override async Task When() - { - await _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(Guid.NewGuid(), "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), new byte[0])); - await _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(Guid.NewGuid(), "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), new byte[0])); - await _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(_id, "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), new byte[0])); - } - - private Task HandleEvent(EventStorePersistentSubscriptionBase sub, ResolvedEvent resolvedEvent) - { - if (_set) - { - return Task.CompletedTask; - } - - _set = true; - _firstEvent = resolvedEvent; - _resetEvent.Set(); - return Task.CompletedTask; - } - - [Test] - public void the_subscription_gets_event_two_as_its_first_event() - { - Assert.IsTrue(_resetEvent.WaitOne(TimeSpan.FromSeconds(10))); - Assert.AreEqual(2, _firstEvent.Event.EventNumber); - Assert.AreEqual(_id, _firstEvent.Event.EventId); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class - connect_to_existing_persistent_subscription_with_start_from_beginning_and_events_in_it_async : - SpecificationWithMiniNode -{ - private readonly string _stream = "$" + Guid.NewGuid(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromBeginning(); - - private readonly AutoResetEvent _resetEvent = new AutoResetEvent(false); - private ResolvedEvent _firstEvent; - private List _ids = new List(); - private bool _set = false; - - private const string _group = "startinbeginning1"; - - protected override async Task Given() - { - await WriteEvents(_conn); - await _conn.CreatePersistentSubscriptionAsync(_stream, _group, _settings, - DefaultData.AdminCredentials); - } - - private async Task WriteEvents(IEventStoreConnection connection) - { - for (int i = 0; i < 10; i++) - { - _ids.Add(Guid.NewGuid()); - await connection.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(_ids[i], "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), new byte[0])); - } - } - - protected override Task When() - { - return _conn.ConnectToPersistentSubscriptionAsync( - _stream, - _group, - HandleEvent, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials); - } - - private Task HandleEvent(EventStorePersistentSubscriptionBase sub, ResolvedEvent resolvedEvent) - { - if (!_set) - { - _set = true; - _firstEvent = resolvedEvent; - _resetEvent.Set(); - } - - return Task.CompletedTask; - } - - [Test] - public void the_subscription_gets_event_zero_as_its_first_event() - { - Assert.IsTrue(_resetEvent.WaitOne(TimeSpan.FromSeconds(10))); - Assert.AreEqual(0, _firstEvent.Event.EventNumber); - Assert.AreEqual(_ids[0], _firstEvent.Event.EventId); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class - connect_to_existing_persistent_subscription_with_start_from_beginning_not_set_and_events_in_it_async : - SpecificationWithMiniNode -{ - private readonly string _stream = "$" + Guid.NewGuid(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - private readonly AutoResetEvent _resetEvent = new AutoResetEvent(false); - - private const string _group = "startinbeginning1"; - - protected override async Task Given() - { - await WriteEvents(_conn); - await _conn.CreatePersistentSubscriptionAsync(_stream, _group, _settings, - DefaultData.AdminCredentials); - } - - private async Task WriteEvents(IEventStoreConnection connection) - { - for (int i = 0; i < 10; i++) - { - await connection.AppendToStreamAsync(_stream, ExpectedVersion.Any, DefaultData.AdminCredentials, - new EventData(Guid.NewGuid(), "test", true, Encoding.UTF8.GetBytes("{'foo' : 'bar'}"), - new byte[0])); - } - } - - protected override Task When() - { - return _conn.ConnectToPersistentSubscriptionAsync( - _stream, - _group, - HandleEvent, - (sub, reason, ex) => { }, - DefaultData.AdminCredentials); - } - - private Task HandleEvent(EventStorePersistentSubscriptionBase sub, ResolvedEvent resolvedEvent) - { - _resetEvent.Set(); - return Task.CompletedTask; - } - - [Test] - public void the_subscription_gets_no_events() - { - Assert.IsFalse(_resetEvent.WaitOne(TimeSpan.FromSeconds(1))); - } -} diff --git a/src/EventStore.Core.Tests/ClientAPI/create_persistent_subscription.cs b/src/EventStore.Core.Tests/ClientAPI/create_persistent_subscription.cs deleted file mode 100644 index 6be3bd0be7..0000000000 --- a/src/EventStore.Core.Tests/ClientAPI/create_persistent_subscription.cs +++ /dev/null @@ -1,305 +0,0 @@ -using System; -using System.Text; -using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.ClientAPI.Exceptions; -using NUnit.Framework; - -namespace EventStore.Core.Tests.ClientAPI; - -[Category("ClientAPI"), Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class create_persistent_subscription_on_existing_stream : SpecificationWithMiniNode -{ - private readonly string _stream = Guid.NewGuid().ToString(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - protected override Task When() - { - return _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, - new EventData(Guid.NewGuid(), "whatever", true, Encoding.UTF8.GetBytes("{'foo' : 2}"), new Byte[0])); - } - - [Test] - public async Task the_completion_succeeds() - { - await _conn.CreatePersistentSubscriptionAsync(_stream, "existing", _settings, - DefaultData.AdminCredentials); - } -} - - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class create_persistent_subscription_on_non_existing_stream : SpecificationWithMiniNode -{ - private readonly string _stream = Guid.NewGuid().ToString(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - protected override Task When() => Task.CompletedTask; - - [Test] - public async Task the_completion_succeeds() - { - await _conn.CreatePersistentSubscriptionAsync(_stream, "nonexistinggroup", _settings, - DefaultData.AdminCredentials); - } -} - - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class create_persistent_subscription_on_all_stream : SpecificationWithMiniNode -{ - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - protected override Task When() => Task.CompletedTask; - - [Test] - public async Task the_completion_fails_with_invalid_stream() - { - await AssertEx.ThrowsAsync(() => - _conn.CreatePersistentSubscriptionAsync("$all", "shitbird", _settings, DefaultData.AdminCredentials)); - } -} - - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class create_persistent_subscription_with_too_big_message_timeout : SpecificationWithMiniNode -{ - protected override Task When() => Task.CompletedTask; - - [Test] - public void the_build_fails_with_argument_exception() - { - Assert.Throws(() => - PersistentSubscriptionSettings.Create().WithMessageTimeoutOf(TimeSpan.FromDays(25 * 365)).Build()); - } -} - - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class create_persistent_subscription_with_too_big_checkpoint_after : SpecificationWithMiniNode -{ - protected override Task When() => Task.CompletedTask; - - [Test] - public void the_build_fails_with_argument_exception() - { - Assert.Throws(() => - PersistentSubscriptionSettings.Create().CheckPointAfter(TimeSpan.FromDays(25 * 365)).Build()); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class create_persistent_subscription_with_dont_timeout : SpecificationWithMiniNode -{ - private readonly string _stream = Guid.NewGuid().ToString(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent() - .DontTimeoutMessages(); - - protected override Task When() => Task.CompletedTask; - - [Test] - public void the_message_timeout_should_be_zero() - { - Assert.That(_settings.MessageTimeout == TimeSpan.Zero); - } - - [Test] - public async Task the_subscription_is_created_without_error() - { - await _conn.CreatePersistentSubscriptionAsync(_stream, "dont-timeout", _settings, - DefaultData.AdminCredentials); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class create_duplicate_persistent_subscription_group : SpecificationWithMiniNode -{ - private readonly string _stream = Guid.NewGuid().ToString(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - protected override Task When() - { - return _conn.CreatePersistentSubscriptionAsync(_stream, "group32", _settings, DefaultData.AdminCredentials); - } - - [Test] - public async Task the_completion_fails_with_invalid_operation_exception() - { - await AssertEx.ThrowsAsync( - () => _conn.CreatePersistentSubscriptionAsync(_stream, "group32", _settings, - DefaultData.AdminCredentials)); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class - can_create_duplicate_persistent_subscription_group_name_on_different_streams - : SpecificationWithMiniNode -{ - private readonly string _stream = Guid.NewGuid().ToString(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - protected override Task When() - { - return _conn.CreatePersistentSubscriptionAsync(_stream, "group3211", _settings, DefaultData.AdminCredentials); - } - - [Test] - public async Task the_completion_succeeds() - { - await - _conn.CreatePersistentSubscriptionAsync("someother" + _stream, "group3211", _settings, - DefaultData.AdminCredentials); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class create_persistent_subscription_group_without_permissions : SpecificationWithMiniNode -{ - private readonly string _stream = Guid.NewGuid().ToString(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - protected override Task When() => Task.CompletedTask; - - [Test] - public async Task the_completion_succeeds() - { - await AssertEx.ThrowsAsync(() => - _conn.CreatePersistentSubscriptionAsync(_stream, "group57", _settings, null)); - } -} - - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class create_persistent_subscription_after_deleting_the_same : SpecificationWithMiniNode -{ - private readonly string _stream = Guid.NewGuid().ToString(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - protected override async Task When() - { - await _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, - new EventData(Guid.NewGuid(), "whatever", true, Encoding.UTF8.GetBytes("{'foo' : 2}"), new Byte[0])); - await _conn.CreatePersistentSubscriptionAsync(_stream, "existing", _settings, DefaultData.AdminCredentials); - await _conn.DeletePersistentSubscriptionAsync(_stream, "existing", DefaultData.AdminCredentials); - } - - [Test] - public async Task the_completion_succeeds() - { - await _conn.CreatePersistentSubscriptionAsync(_stream, "existing", _settings, DefaultData.AdminCredentials); - } -} - -//ALL -/* - - [TestFixture, Category("LongRunning")] - public class create_persistent_subscription_on_all : SpecificationWithMiniNode - { - private PersistentSubscriptionCreateResult _result; - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettingsBuilder.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - protected override void When() - { - _result = _conn.CreatePersistentSubscriptionForAllAsync("group", _settings, DefaultData.AdminCredentials).Result; - } - - [Test] - public void the_completion_succeeds() - { - Assert.AreEqual(PersistentSubscriptionCreateStatus.Success, _result.Status); - } - } - - - [TestFixture, Category("LongRunning")] - public class create_duplicate_persistent_subscription_group_on_all : SpecificationWithMiniNode - { - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettingsBuilder.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - protected override void When() - { - _conn.CreatePersistentSubscriptionForAllAsync("group32", _settings, DefaultData.AdminCredentials).Wait(); - } - - [Test] - public void the_completion_fails_with_invalid_operation_exception() - { - try - { - _conn.CreatePersistentSubscriptionForAllAsync("group32", _settings, DefaultData.AdminCredentials).Wait(); - throw new Exception("expected exception"); - } - catch (Exception ex) - { - Assert.IsInstanceOf(typeof(AggregateException), ex); - var inner = ex.InnerException; - Assert.IsInstanceOf(typeof(InvalidOperationException), inner); - } - } - } - - [TestFixture, Category("LongRunning")] - public class create_persistent_subscription_group_on_all_without_permissions : SpecificationWithMiniNode - { - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettingsBuilder.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - protected override void When() - { - } - - [Test] - public void the_completion_succeeds() - { - try - { - _conn.CreatePersistentSubscriptionForAllAsync("group57", _settings, null).Wait(); - throw new Exception("expected exception"); - } - catch (Exception ex) - { - Assert.IsInstanceOf(typeof(AggregateException), ex); - var inner = ex.InnerException; - Assert.IsInstanceOf(typeof(AccessDeniedException), inner); - } - } - } -*/ diff --git a/src/EventStore.Core.Tests/ClientAPI/deleting_persistent_subscription.cs b/src/EventStore.Core.Tests/ClientAPI/deleting_persistent_subscription.cs deleted file mode 100644 index 641589ed1a..0000000000 --- a/src/EventStore.Core.Tests/ClientAPI/deleting_persistent_subscription.cs +++ /dev/null @@ -1,173 +0,0 @@ -using System; -using System.Threading; -using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.ClientAPI.Exceptions; -using EventStore.Core.Tests.Helpers; -using NUnit.Framework; - -namespace EventStore.Core.Tests.ClientAPI; - -[Category("ClientAPI"), Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class deleting_existing_persistent_subscription_group_with_permissions - : SpecificationWithMiniNode -{ - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - private readonly string _stream = Guid.NewGuid().ToString(); - - protected override Task When() => - _conn.CreatePersistentSubscriptionAsync(_stream, "groupname123", _settings, - DefaultData.AdminCredentials); - - [Test] - public async Task the_delete_of_group_succeeds() - { - await _conn.DeletePersistentSubscriptionAsync(_stream, "groupname123", DefaultData.AdminCredentials); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class deleting_existing_persistent_subscription_with_subscriber : SpecificationWithMiniNode -{ - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - private readonly string _stream = Guid.NewGuid().ToString(); - private readonly ManualResetEvent _called = new ManualResetEvent(false); - - protected override async Task Given() - { - await base.Given(); - await _conn.CreatePersistentSubscriptionAsync(_stream, "groupname123", _settings, - DefaultData.AdminCredentials); - _conn.ConnectToPersistentSubscription(_stream, "groupname123", - (s, e) => Task.CompletedTask, - (s, r, e) => _called.Set(), DefaultData.AdminCredentials); - } - - protected override Task When() - { - return _conn.DeletePersistentSubscriptionAsync(_stream, "groupname123", DefaultData.AdminCredentials); - } - - [Test] - public void the_subscription_is_dropped() - { - Assert.IsTrue(_called.WaitOne(TimeSpan.FromSeconds(5))); - } -} - - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class deleting_persistent_subscription_group_that_doesnt_exist : SpecificationWithMiniNode -{ - private readonly string _stream = Guid.NewGuid().ToString(); - - protected override Task When() => Task.CompletedTask; - - [Test] - public async Task the_delete_fails_with_argument_exception() - { - await AssertEx.ThrowsAsync( - () => - _conn.DeletePersistentSubscriptionAsync(_stream, Guid.NewGuid().ToString(), - DefaultData.AdminCredentials)); - } -} - - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class deleting_persistent_subscription_group_without_permissions : SpecificationWithMiniNode -{ - private readonly string _stream = Guid.NewGuid().ToString(); - - protected override Task When() => Task.CompletedTask; - - [Test] - public async Task the_delete_fails_with_access_denied() - { - await AssertEx.ThrowsAsync( - () => _conn.DeletePersistentSubscriptionAsync(_stream, Guid.NewGuid().ToString())); - } -} - -//ALL -/* - - [TestFixture, Category("LongRunning")] - public class deleting_existing_persistent_subscription_group_on_all_with_permissions : SpecificationWithMiniNode - { - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettingsBuilder.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - protected override void When() - { - _conn.CreatePersistentSubscriptionForAllAsync("groupname123", _settings, - DefaultData.AdminCredentials).Wait(); - } - - [Test] - public void the_delete_of_group_succeeds() - { - var result = _conn.DeletePersistentSubscriptionForAllAsync("groupname123", DefaultData.AdminCredentials).Result; - Assert.AreEqual(PersistentSubscriptionDeleteStatus.Success, result.Status); - } - } - - [TestFixture, Category("LongRunning")] - public class deleting_persistent_subscription_group_on_all_that_doesnt_exist : SpecificationWithMiniNode - { - protected override void When() - { - } - - [Test] - public void the_delete_fails_with_argument_exception() - { - try - { - _conn.DeletePersistentSubscriptionForAllAsync(Guid.NewGuid().ToString(), DefaultData.AdminCredentials).Wait(); - throw new Exception("expected exception"); - } - catch (Exception ex) - { - Assert.IsInstanceOf(typeof(AggregateException), ex); - var inner = ex.InnerException; - Assert.IsInstanceOf(typeof(InvalidOperationException), inner); - } - } - } - - - [TestFixture, Category("LongRunning")] - public class deleting_persistent_subscription_group_on_all_without_permissions : SpecificationWithMiniNode - { - protected override void When() - { - } - - [Test] - public void the_delete_fails_with_access_denied() - { - try - { - _conn.DeletePersistentSubscriptionForAllAsync(Guid.NewGuid().ToString()).Wait(); - throw new Exception("expected exception"); - } - catch (Exception ex) - { - Assert.IsInstanceOf(typeof(AggregateException), ex); - var inner = ex.InnerException; - Assert.IsInstanceOf(typeof(AccessDeniedException), inner); - } - } - } -*/ diff --git a/src/EventStore.Core.Tests/ClientAPI/persistent_connect_integration_tests.cs b/src/EventStore.Core.Tests/ClientAPI/persistent_connect_integration_tests.cs deleted file mode 100644 index 9479e4c9c0..0000000000 --- a/src/EventStore.Core.Tests/ClientAPI/persistent_connect_integration_tests.cs +++ /dev/null @@ -1,487 +0,0 @@ -using System; -using System.Text; -using System.Threading; -using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.ClientAPI.Common; -using NUnit.Framework; - -namespace EventStore.Core.Tests.ClientAPI; - -[Category("ClientAPI"), Category("LongRunning")] -[NonParallelizable] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class happy_case_writing_and_subscribing_to_normal_events_manual_ack - : SpecificationWithMiniNode -{ - private const int BufferCount = 10; - private const int EventWriteCount = BufferCount * 2; - - private readonly ManualResetEvent _eventsReceived = new ManualResetEvent(false); - private int _eventReceivedCount; - - protected override Task When() => Task.CompletedTask; - - [Test] - public async Task Test() - { - var streamName = Guid.NewGuid().ToString(); - var groupName = Guid.NewGuid().ToString(); - var settings = PersistentSubscriptionSettings - .Create() - .StartFromCurrent() - .ResolveLinkTos() - .Build(); - - await _conn.CreatePersistentSubscriptionAsync(streamName, groupName, settings, DefaultData.AdminCredentials); - await _conn.ConnectToPersistentSubscriptionAsync(streamName, groupName, - (subscription, resolvedEvent) => - { - subscription.Acknowledge(resolvedEvent); - - if (Interlocked.Increment(ref _eventReceivedCount) == EventWriteCount) - { - _eventsReceived.Set(); - } - - return Task.CompletedTask; - }, - (sub, reason, exception) => - Console.WriteLine("Subscription dropped (reason:{0}, exception:{1}).", reason, exception), - bufferSize: 10, autoAck: false, userCredentials: DefaultData.AdminCredentials); - - for (var i = 0; i < EventWriteCount; i++) - { - var eventData = new EventData(Guid.NewGuid(), "SomeEvent", false, new byte[0], new byte[0]); - - await _conn.AppendToStreamAsync(streamName, ExpectedVersion.Any, DefaultData.AdminCredentials, eventData); - } - - if (!_eventsReceived.WaitOne(TimeSpan.FromSeconds(5))) - { - throw new Exception("Timed out waiting for events."); - } - } -} - - -[Category("LongRunning")] -[NonParallelizable] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class happy_case_writing_and_subscribing_to_normal_events_auto_ack : SpecificationWithMiniNode -{ - private const int BufferCount = 10; - private const int EventWriteCount = BufferCount * 2; - - private readonly ManualResetEvent _eventsReceived = new ManualResetEvent(false); - private int _eventReceivedCount; - - protected override Task When() => Task.CompletedTask; - - [Test] - public async Task Test() - { - var streamName = Guid.NewGuid().ToString(); - var groupName = Guid.NewGuid().ToString(); - var settings = PersistentSubscriptionSettings - .Create() - .StartFromCurrent() - .ResolveLinkTos() - .Build(); - await _conn.CreatePersistentSubscriptionAsync(streamName, groupName, settings, DefaultData.AdminCredentials) -; - await _conn.ConnectToPersistentSubscriptionAsync(streamName, groupName, - (subscription, resolvedEvent) => - { - if (Interlocked.Increment(ref _eventReceivedCount) == EventWriteCount) - { - _eventsReceived.Set(); - } - - return Task.CompletedTask; - }, - (sub, reason, exception) => - Console.WriteLine("Subscription dropped (reason:{0}, exception:{1}).", reason, exception), - userCredentials: DefaultData.AdminCredentials); - - for (var i = 0; i < EventWriteCount; i++) - { - var eventData = new EventData(Guid.NewGuid(), "SomeEvent", false, new byte[0], new byte[0]); - - await _conn.AppendToStreamAsync(streamName, ExpectedVersion.Any, DefaultData.AdminCredentials, eventData); - } - - if (!_eventsReceived.WaitOne(TimeSpan.FromSeconds(15))) - { - throw new Exception($"Timed out waiting for events, received {_eventReceivedCount} event(s)."); - } - } -} - - -[Category("LongRunning")] -[NonParallelizable] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class happy_case_catching_up_to_normal_events_auto_ack : SpecificationWithMiniNode -{ - private const int BufferCount = 10; - private const int EventWriteCount = BufferCount * 2; - - private readonly ManualResetEvent _eventsReceived = new ManualResetEvent(false); - private int _eventReceivedCount; - - protected override Task When() => Task.CompletedTask; - - - [Test] - public async Task Test() - { - var streamName = Guid.NewGuid().ToString(); - var groupName = Guid.NewGuid().ToString(); - var settings = PersistentSubscriptionSettings - .Create() - .StartFromBeginning() - .ResolveLinkTos() - .Build(); - for (var i = 0; i < EventWriteCount; i++) - { - var eventData = new EventData(Guid.NewGuid(), "SomeEvent", false, new byte[0], new byte[0]); - - await _conn.AppendToStreamAsync(streamName, ExpectedVersion.Any, DefaultData.AdminCredentials, eventData); - } - - await _conn.CreatePersistentSubscriptionAsync(streamName, groupName, settings, DefaultData.AdminCredentials); - await _conn.ConnectToPersistentSubscriptionAsync(streamName, groupName, - (subscription, resolvedEvent) => - { - if (Interlocked.Increment(ref _eventReceivedCount) == EventWriteCount) - { - _eventsReceived.Set(); - } - - return Task.CompletedTask; - }, - (sub, reason, exception) => - { - Console.WriteLine($"Subscription dropped (reason:{reason}, exception:{exception})."); - }, - userCredentials: DefaultData.AdminCredentials, - autoAck: true); - - if (!_eventsReceived.WaitOne(TimeSpan.FromSeconds(5))) - { - throw new Exception("Timed out waiting for events."); - } - } -} - - -[Category("LongRunning")] -[NonParallelizable] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class happy_case_catching_up_to_normal_events_manual_ack : SpecificationWithMiniNode -{ - private const int BufferCount = 10; - private const int EventWriteCount = BufferCount * 2; - - private readonly ManualResetEvent _eventsReceived = new ManualResetEvent(false); - private int _eventReceivedCount; - - protected override Task When() => Task.CompletedTask; - - - [Test] - public async Task Test() - { - var streamName = Guid.NewGuid().ToString(); - var groupName = Guid.NewGuid().ToString(); - var settings = PersistentSubscriptionSettings - .Create() - .StartFromBeginning() - .ResolveLinkTos() - .Build(); - for (var i = 0; i < EventWriteCount; i++) - { - var eventData = new EventData(Guid.NewGuid(), "SomeEvent", false, new byte[0], new byte[0]); - - await _conn.AppendToStreamAsync(streamName, ExpectedVersion.Any, DefaultData.AdminCredentials, eventData); - } - - await _conn.CreatePersistentSubscriptionAsync(streamName, groupName, settings, DefaultData.AdminCredentials); - await _conn.ConnectToPersistentSubscriptionAsync(streamName, groupName, - (subscription, resolvedEvent) => - { - subscription.Acknowledge(resolvedEvent); - if (Interlocked.Increment(ref _eventReceivedCount) == EventWriteCount) - { - _eventsReceived.Set(); - } - - return Task.CompletedTask; - }, - (sub, reason, exception) => - Console.WriteLine("Subscription dropped (reason:{0}, exception:{1}).", reason, exception), - userCredentials: DefaultData.AdminCredentials, - autoAck: false); - - if (!_eventsReceived.WaitOne(TimeSpan.FromSeconds(5))) - { - throw new Exception("Timed out waiting for events."); - } - } -} - - -[Category("LongRunning")] -[NonParallelizable] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class happy_case_catching_up_to_link_to_events_manual_ack : SpecificationWithMiniNode -{ - private const int BufferCount = 10; - private const int EventWriteCount = BufferCount * 2; - - private readonly ManualResetEvent _eventsReceived = new ManualResetEvent(false); - private int _eventReceivedCount; - - protected override Task When() => Task.CompletedTask; - - [Test] - public async Task Test() - { - var streamName = Guid.NewGuid().ToString(); - var groupName = Guid.NewGuid().ToString(); - var settings = PersistentSubscriptionSettings - .Create() - .StartFromBeginning() - .ResolveLinkTos() - .Build(); - for (var i = 0; i < EventWriteCount; i++) - { - var eventData = new EventData(Guid.NewGuid(), "SomeEvent", false, new byte[0], new byte[0]); - await _conn.AppendToStreamAsync(streamName + "original", ExpectedVersion.Any, DefaultData.AdminCredentials, - eventData); - } - - await _conn.CreatePersistentSubscriptionAsync(streamName, groupName, settings, DefaultData.AdminCredentials); - await _conn.ConnectToPersistentSubscriptionAsync(streamName, groupName, - (subscription, resolvedEvent) => - { - subscription.Acknowledge(resolvedEvent); - if (Interlocked.Increment(ref _eventReceivedCount) == EventWriteCount) - { - _eventsReceived.Set(); - } - - return Task.CompletedTask; - }, - (sub, reason, exception) => - Console.WriteLine("Subscription dropped (reason:{0}, exception:{1}).", reason, exception), - userCredentials: DefaultData.AdminCredentials, - autoAck: false); - for (var i = 0; i < EventWriteCount; i++) - { - var eventData = new EventData(Guid.NewGuid(), SystemEventTypes.LinkTo, false, - Encoding.UTF8.GetBytes(i + "@" + streamName + "original"), null); - await _conn.AppendToStreamAsync(streamName, ExpectedVersion.Any, DefaultData.AdminCredentials, eventData); - } - - if (!_eventsReceived.WaitOne(TimeSpan.FromSeconds(5))) - { - throw new Exception("Timed out waiting for events."); - } - } -} - - -[Category("LongRunning")] -[NonParallelizable] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class happy_case_catching_up_to_link_to_events_auto_ack : SpecificationWithMiniNode -{ - private const int BufferCount = 10; - private const int EventWriteCount = BufferCount * 2; - - private readonly ManualResetEvent _eventsReceived = new ManualResetEvent(false); - private int _eventReceivedCount; - - protected override Task When() => Task.CompletedTask; - - [Test] - public async Task Test() - { - var streamName = Guid.NewGuid().ToString(); - var groupName = Guid.NewGuid().ToString(); - var settings = PersistentSubscriptionSettings - .Create() - .StartFromBeginning() - .ResolveLinkTos() - .Build(); - for (var i = 0; i < EventWriteCount; i++) - { - var eventData = new EventData(Guid.NewGuid(), "SomeEvent", false, new byte[0], new byte[0]); - await _conn.AppendToStreamAsync(streamName + "original", ExpectedVersion.Any, DefaultData.AdminCredentials, - eventData); - } - - await _conn.CreatePersistentSubscriptionAsync(streamName, groupName, settings, DefaultData.AdminCredentials); - await _conn.ConnectToPersistentSubscriptionAsync(streamName, groupName, - (subscription, resolvedEvent) => - { - if (Interlocked.Increment(ref _eventReceivedCount) == EventWriteCount) - { - _eventsReceived.Set(); - } - - return Task.CompletedTask; - }, - (sub, reason, exception) => - Console.WriteLine("Subscription dropped (reason:{0}, exception:{1}).", reason, exception), - userCredentials: DefaultData.AdminCredentials, - autoAck: true); - for (var i = 0; i < EventWriteCount; i++) - { - var eventData = new EventData(Guid.NewGuid(), SystemEventTypes.LinkTo, false, - Encoding.UTF8.GetBytes(i + "@" + streamName + "original"), null); - await _conn.AppendToStreamAsync(streamName, ExpectedVersion.Any, DefaultData.AdminCredentials, eventData); - } - - if (!_eventsReceived.WaitOne(TimeSpan.FromSeconds(5))) - { - throw new Exception("Timed out waiting for events."); - } - } -} - -[Category("ClientAPI"), Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class when_writing_and_subscribing_to_normal_events_manual_nack - : SpecificationWithMiniNode -{ - private const int BufferCount = 10; - private const int EventWriteCount = BufferCount * 2; - - private readonly ManualResetEvent _eventsReceived = new ManualResetEvent(false); - private int _eventReceivedCount; - - protected override Task When() => Task.CompletedTask; - - [Test] - public async Task Test() - { - var streamName = Guid.NewGuid().ToString(); - var groupName = Guid.NewGuid().ToString(); - var settings = PersistentSubscriptionSettings - .Create() - .StartFromCurrent() - .ResolveLinkTos() - .Build(); - - await _conn.CreatePersistentSubscriptionAsync(streamName, groupName, settings, DefaultData.AdminCredentials); - await _conn.ConnectToPersistentSubscriptionAsync(streamName, groupName, - (subscription, resolvedEvent) => - { - subscription.Fail(resolvedEvent, PersistentSubscriptionNakEventAction.Park, "fail"); - - if (Interlocked.Increment(ref _eventReceivedCount) == EventWriteCount) - { - _eventsReceived.Set(); - } - - return Task.CompletedTask; - }, - (sub, reason, exception) => - Console.WriteLine("Subscription dropped (reason:{0}, exception:{1}).", reason, exception), - bufferSize: 10, autoAck: false, userCredentials: DefaultData.AdminCredentials); - - for (var i = 0; i < EventWriteCount; i++) - { - var eventData = new EventData(Guid.NewGuid(), "SomeEvent", false, new byte[0], new byte[0]); - - await _conn.AppendToStreamAsync(streamName, ExpectedVersion.Any, DefaultData.AdminCredentials, eventData); - } - - if (!_eventsReceived.WaitOne(TimeSpan.FromSeconds(5))) - { - throw new Exception("Timed out waiting for events."); - } - } -} - -[Category("ClientAPI"), Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class when_connection_drops_messages_that_have_run_out_of_retries_are_not_retried - : SpecificationWithMiniNode -{ - private readonly TaskCompletionSource _subscriptionDropped = new TaskCompletionSource(); - private readonly TaskCompletionSource _eventReceived = new TaskCompletionSource(); - private ResolvedEvent _receivedEvent; - - protected override Task When() => Task.CompletedTask; - - [Test] - [Retry(10)] - public void Test() => Assert.DoesNotThrowAsync(async () => - { - await CloseConnectionAndWait(_conn); - _conn = BuildConnection(_node); - AddLogging(_conn); - await _conn.ConnectAsync(); - - var streamName = Guid.NewGuid().ToString(); - var groupName = Guid.NewGuid().ToString(); - var settings = PersistentSubscriptionSettings - .Create() - .StartFromCurrent() - .WithMaxRetriesOf(0) // Don't retry messages - .Build(); - - await _conn.CreatePersistentSubscriptionAsync(streamName, groupName, settings, DefaultData.AdminCredentials); - await _conn.ConnectToPersistentSubscriptionAsync(streamName, groupName, async (subscription, resolvedEvent) => - { - await CloseConnectionAndWait(_conn); - }, - (sub, reason, exception) => - { - Console.WriteLine("Subscription dropped (reason:{0}, exception:{1}). @ {2}", reason, exception, DateTime.Now); - _subscriptionDropped.TrySetResult(true); - }, - bufferSize: 10, autoAck: false, userCredentials: DefaultData.AdminCredentials); - - var parkedEventData = new EventData(Guid.NewGuid(), "SomeEvent", false, new byte[0], new byte[0]); - await _conn.AppendToStreamAsync(streamName, ExpectedVersion.Any, DefaultData.AdminCredentials, parkedEventData); - - await _subscriptionDropped.Task.WithTimeout(); - - _conn = BuildConnection(_node); - AddLogging(_conn); - await _conn.ConnectAsync(); - - await _conn.ConnectToPersistentSubscriptionAsync(streamName, groupName, (subscription, resolvedEvent) => - { - subscription.Acknowledge(resolvedEvent); - _receivedEvent = resolvedEvent; - _eventReceived.TrySetResult(true); - }, (sub, reason, exception) => - { - Console.WriteLine("Second Subscription dropped (reason:{0}, exception:{1}).", reason, exception); - }, bufferSize: 10, autoAck: false, userCredentials: DefaultData.AdminCredentials); - - // Ensure we only get the new event, not the previous one - var newEventData = new EventData(Guid.NewGuid(), "SomeEvent", false, new byte[0], new byte[0]); - await _conn.AppendToStreamAsync(streamName, ExpectedVersion.Any, DefaultData.AdminCredentials, newEventData); - - await _eventReceived.Task.WithTimeout(); - Assert.AreEqual(newEventData.EventId, _receivedEvent.Event.EventId); - - //flaky: temporarily added for debugging - void AddLogging(IEventStoreConnection conn) - { - conn.AuthenticationFailed += (_, args) => Console.WriteLine($"_conn.AuthenticationFailed: {args.Connection.ConnectionName} @ {DateTime.Now} {TestContext.CurrentContext.CurrentRepeatCount}"); - conn.Closed += (_, args) => Console.WriteLine($"_conn.Closed: {args.Connection.ConnectionName} @ {DateTime.Now} {TestContext.CurrentContext.CurrentRepeatCount}"); - conn.Connected += (_, args) => Console.WriteLine($"_conn.Connected: {args.Connection.ConnectionName} @ {DateTime.Now} {TestContext.CurrentContext.CurrentRepeatCount}"); - conn.Disconnected += (_, args) => Console.WriteLine($"_conn.Disconnected: {args.Connection.ConnectionName} @ {DateTime.Now} {TestContext.CurrentContext.CurrentRepeatCount}"); - conn.ErrorOccurred += (_, args) => Console.WriteLine($"_conn.ErrorOccurred: {args.Connection.ConnectionName} @ {DateTime.Now} {TestContext.CurrentContext.CurrentRepeatCount}"); - conn.Reconnecting += (_, args) => Console.WriteLine($"_conn.Reconnecting: {args.Connection.ConnectionName} @ {DateTime.Now} {TestContext.CurrentContext.CurrentRepeatCount}"); - } - }); -} diff --git a/src/EventStore.Core.Tests/ClientAPI/read_from_persistent_subscription_with_link_resolution_when_stream_name_contains_at_symbol.cs b/src/EventStore.Core.Tests/ClientAPI/read_from_persistent_subscription_with_link_resolution_when_stream_name_contains_at_symbol.cs deleted file mode 100644 index 014e6807c0..0000000000 --- a/src/EventStore.Core.Tests/ClientAPI/read_from_persistent_subscription_with_link_resolution_when_stream_name_contains_at_symbol.cs +++ /dev/null @@ -1,48 +0,0 @@ -using System; -using System.Text; -using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.Core.Tests.ClientAPI.Helpers; -using NUnit.Framework; -using ExpectedVersion = EventStore.ClientAPI.ExpectedVersion; - -namespace EventStore.Core.Tests.ClientAPI; - -[Category("LongRunning"), Category("ClientAPI")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class read_from_persistent_subscription_with_link_resolution_when_stream_name_contains_at_symbol : SpecificationWithMiniNode -{ - private string _result; - - protected override async Task When() - { - var task = new TaskCompletionSource(); - - var setts = PersistentSubscriptionSettings.Create() - .ResolveLinkTos() - .StartFromBeginning(); - - await _conn.CreatePersistentSubscriptionAsync("link", "Agroup", setts, DefaultData.AdminCredentials); - await _conn.ConnectToPersistentSubscriptionAsync( - "link", - "Agroup", - (sub, @event) => - { - var data = Encoding.Default.GetString(@event.Event.Data); - task.TrySetResult(data); - return Task.CompletedTask; - }, - (sub, reason, ex) => { }, DefaultData.AdminCredentials); - - await _conn.AppendToStreamAsync("target@me@com", ExpectedVersion.NoStream, TestEvent.NewTestEvent("data", eventName: "AEvent")); - await _conn.AppendToStreamAsync("link", ExpectedVersion.NoStream, TestEvent.NewTestEvent("0@target@me@com", eventName: "$>")); - - _result = await Task.WhenAny(task.Task, Task.Delay(TimeSpan.FromSeconds(30)).ContinueWith(_ => "timeout")).Result; - } - - [Test] - public void the_subscription_resolve_the_link_properly() - { - Assert.AreEqual(_result, "data"); - } -} diff --git a/src/EventStore.Core.Tests/ClientAPI/update_persistent_subscription.cs b/src/EventStore.Core.Tests/ClientAPI/update_persistent_subscription.cs deleted file mode 100644 index 087e884922..0000000000 --- a/src/EventStore.Core.Tests/ClientAPI/update_persistent_subscription.cs +++ /dev/null @@ -1,140 +0,0 @@ -using System; -using System.Text; -using System.Threading; -using System.Threading.Tasks; -using EventStore.ClientAPI; -using EventStore.ClientAPI.Exceptions; -using NUnit.Framework; - -namespace EventStore.Core.Tests.ClientAPI; - -[Category("ClientAPI"), Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class update_existing_persistent_subscription : SpecificationWithMiniNode -{ - private readonly string _stream = Guid.NewGuid().ToString(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - protected override async Task Given() - { - await _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, - new EventData(Guid.NewGuid(), "whatever", true, Encoding.UTF8.GetBytes("{'foo' : 2}"), new Byte[0])); - await _conn.CreatePersistentSubscriptionAsync(_stream, "existing", _settings, DefaultData.AdminCredentials); - } - - protected override Task When() => Task.CompletedTask; - - [Test] - public async Task the_completion_succeeds() - { - await _conn.UpdatePersistentSubscriptionAsync(_stream, "existing", _settings, DefaultData.AdminCredentials); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class update_existing_persistent_subscription_with_subscribers : SpecificationWithMiniNode -{ - private readonly string _stream = Guid.NewGuid().ToString(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - private readonly AutoResetEvent _dropped = new AutoResetEvent(false); - private SubscriptionDropReason _reason; - private Exception _exception; - private Exception _caught = null; - - protected override async Task Given() - { - await _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, - new EventData(Guid.NewGuid(), "whatever", true, Encoding.UTF8.GetBytes("{'foo' : 2}"), new Byte[0])); - await _conn.CreatePersistentSubscriptionAsync(_stream, "existing", _settings, DefaultData.AdminCredentials) -; - _conn.ConnectToPersistentSubscription(_stream, "existing", (x, y) => Task.CompletedTask, - (sub, reason, ex) => - { - _dropped.Set(); - _reason = reason; - _exception = ex; - }, DefaultData.AdminCredentials); - } - - protected override async Task When() - { - try - { - await _conn.UpdatePersistentSubscriptionAsync(_stream, "existing", _settings, DefaultData.AdminCredentials); - } - catch (Exception ex) - { - _caught = ex; - } - } - - [Test] - public void the_completion_succeeds() - { - Assert.IsNull(_caught); - } - - [Test] - public void existing_subscriptions_are_dropped() - { - Assert.IsTrue(_dropped.WaitOne(TimeSpan.FromSeconds(5))); - Assert.AreEqual(SubscriptionDropReason.UserInitiated, _reason); - Assert.IsNull(_exception); - } -} - - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class update_non_existing_persistent_subscription : SpecificationWithMiniNode -{ - private readonly string _stream = Guid.NewGuid().ToString(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - protected override Task When() => Task.CompletedTask; - - [Test] - public async Task the_completion_fails_with_not_found() - { - await AssertEx.ThrowsAsync( - () => _conn.UpdatePersistentSubscriptionAsync(_stream, "existing", _settings, - DefaultData.AdminCredentials)); - } -} - -[Category("LongRunning")] -[TestFixture(typeof(LogFormat.V2), typeof(string))] -public class update_existing_persistent_subscription_without_permissions : SpecificationWithMiniNode -{ - private readonly string _stream = Guid.NewGuid().ToString(); - - private readonly PersistentSubscriptionSettings _settings = PersistentSubscriptionSettings.Create() - .DoNotResolveLinkTos() - .StartFromCurrent(); - - protected override async Task When() - { - await _conn.AppendToStreamAsync(_stream, ExpectedVersion.Any, - new EventData(Guid.NewGuid(), "whatever", true, Encoding.UTF8.GetBytes("{'foo' : 2}"), new Byte[0])); - await _conn.CreatePersistentSubscriptionAsync(_stream, "existing", _settings, DefaultData.AdminCredentials) -; - } - - [Test] - public async Task the_completion_fails_with_access_denied() - { - await AssertEx.ThrowsAsync( - () => _conn.UpdatePersistentSubscriptionAsync(_stream, "existing", _settings, null)); - } -} diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/PersistentSubscriptionTests/CreateTests.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/PersistentSubscriptionTests/CreateTests.cs index 2e2dbef364..77f24ff488 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Grpc/PersistentSubscriptionTests/CreateTests.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/PersistentSubscriptionTests/CreateTests.cs @@ -31,6 +31,14 @@ public async Task can_create_persistent_subscription() await client.CreateAsync(CreateRequest(), GetCallOptions(AdminCredentials)); } + [Test] + public async Task can_create_persistent_subscription_to_all() + { + var client = new PersistentSubscriptions.PersistentSubscriptionsClient(Channel); + + await client.CreateAsync(CreateAllRequest(), GetCallOptions(AdminCredentials)); + } + [Test] public async Task creating_duplicate_persistent_subscription_returns_already_exists() { @@ -47,6 +55,55 @@ public async Task creating_duplicate_persistent_subscription_returns_already_exi Assert.AreEqual(StatusCode.AlreadyExists, ex.Status.StatusCode); } + [Test] + public async Task creating_duplicate_persistent_subscription_to_all_returns_already_exists() + { + var client = new PersistentSubscriptions.PersistentSubscriptionsClient(Channel); + var request = CreateAllRequest(); + + await client.CreateAsync(request, GetCallOptions(AdminCredentials)); + + var ex = Assert.ThrowsAsync(async () => + await client.CreateAsync(request, GetCallOptions(AdminCredentials))); + + Assert.AreEqual(StatusCode.AlreadyExists, ex.Status.StatusCode); + } + + [Test] + public async Task can_reuse_a_group_name_on_a_different_stream() + { + var client = new PersistentSubscriptions.PersistentSubscriptionsClient(Channel); + var groupName = NewName("group"); + + await client.CreateAsync( + CreateRequest(NewName("stream"), groupName), GetCallOptions(AdminCredentials)); + await client.CreateAsync( + CreateRequest(NewName("stream"), groupName), GetCallOptions(AdminCredentials)); + } + + [Test] + public async Task can_recreate_a_subscription_after_deleting_it() + { + var client = new PersistentSubscriptions.PersistentSubscriptionsClient(Channel); + var streamName = NewName("stream"); + var groupName = NewName("group"); + var request = CreateRequest(streamName, groupName); + await client.CreateAsync(request, GetCallOptions(AdminCredentials)); + await client.DeleteAsync(new DeleteReq + { + Options = new DeleteReq.Types.Options + { + GroupName = groupName, + StreamIdentifier = new StreamIdentifier + { + StreamName = ByteString.CopyFromUtf8(streamName) + } + } + }, GetCallOptions(AdminCredentials)); + + await client.CreateAsync(request, GetCallOptions(AdminCredentials)); + } + [Test] public void creating_persistent_subscription_without_permissions_returns_permission_denied() { @@ -58,6 +115,17 @@ public void creating_persistent_subscription_without_permissions_returns_permiss Assert.AreEqual(StatusCode.PermissionDenied, ex.Status.StatusCode); } + [Test] + public void creating_persistent_subscription_to_all_without_permissions_returns_permission_denied() + { + var client = new PersistentSubscriptions.PersistentSubscriptionsClient(Channel); + + var ex = Assert.ThrowsAsync(async () => + await client.CreateAsync(CreateAllRequest(), GetCallOptions())); + + Assert.AreEqual(StatusCode.PermissionDenied, ex.Status.StatusCode); + } + [Test] public void creating_persistent_subscription_with_bad_config_returns_invalid_argument() { @@ -111,6 +179,20 @@ private CreateReq CreateRequest( } }; + private CreateReq CreateAllRequest(string groupName = null) => new() + { + Options = new CreateReq.Types.Options + { + GroupName = groupName ?? NewName("group"), + All = new CreateReq.Types.AllOptions + { + Start = new Empty(), + NoFilter = new Empty() + }, + Settings = Settings() + } + }; + private CreateReq.Types.Settings Settings(int messageTimeoutMs = 20000, bool includeMessageTimeout = true) { var settings = new CreateReq.Types.Settings diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/PersistentSubscriptionTests/DeleteTests.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/PersistentSubscriptionTests/DeleteTests.cs index 64211dcee0..a19d0f6779 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Grpc/PersistentSubscriptionTests/DeleteTests.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/PersistentSubscriptionTests/DeleteTests.cs @@ -168,6 +168,61 @@ public void returns_not_found() } } + [TestFixture(typeof(LogFormat.V2), typeof(string))] + public class deleting_persistent_subscriptions_to_all + : GrpcSpecification + { + private PersistentSubscriptions.PersistentSubscriptionsClient _client; + + protected override Task Given() + { + _client = new PersistentSubscriptions.PersistentSubscriptionsClient(Channel); + return Task.CompletedTask; + } + + protected override Task When() => Task.CompletedTask; + + [Test] + public async Task removes_an_existing_subscription() + { + var groupName = Guid.NewGuid().ToString(); + await _client.CreateAsync(CreateAllRequest(groupName), GetCallOptions(AdminCredentials)); + + await _client.DeleteAsync(DeleteAllRequest(groupName), GetCallOptions(AdminCredentials)); + + var exception = Assert.ThrowsAsync(async () => + await _client.GetInfoAsync(GetAllInfoRequest(groupName), GetCallOptions(AdminCredentials))); + Assert.AreEqual(StatusCode.NotFound, exception.StatusCode); + } + + [Test] + public void deleting_a_missing_subscription_returns_not_found() + { + var exception = Assert.ThrowsAsync(async () => + await _client.DeleteAsync( + DeleteAllRequest(Guid.NewGuid().ToString()), + GetCallOptions(AdminCredentials))); + + Assert.AreEqual(StatusCode.NotFound, exception.StatusCode); + } + + [Test] + public async Task deleting_without_permission_returns_permission_denied_and_preserves_the_subscription() + { + var groupName = Guid.NewGuid().ToString(); + await _client.CreateAsync(CreateAllRequest(groupName), GetCallOptions(AdminCredentials)); + + var exception = Assert.ThrowsAsync(async () => + await _client.DeleteAsync(DeleteAllRequest(groupName), GetCallOptions())); + Assert.AreEqual(StatusCode.PermissionDenied, exception.StatusCode); + + var response = await _client.GetInfoAsync( + GetAllInfoRequest(groupName), + GetCallOptions(AdminCredentials)); + Assert.AreEqual(groupName, response.SubscriptionInfo.GroupName); + } + } + private static async Task> SubscribeToPersistentSubscription( PersistentSubscriptions.PersistentSubscriptionsClient client, string streamName, string groupName, CallOptions callOptions) { @@ -216,6 +271,20 @@ await call.RequestStream.WriteAsync(new ReadReq } }; + private static CreateReq CreateAllRequest(string groupName) => new() + { + Options = new CreateReq.Types.Options + { + GroupName = groupName, + All = new CreateReq.Types.AllOptions + { + Start = new Empty(), + NoFilter = new Empty() + }, + Settings = Settings + } + }; + private static DeleteReq DeleteRequest(string streamName, string groupName) => new() { Options = new DeleteReq.Types.Options @@ -228,6 +297,15 @@ await call.RequestStream.WriteAsync(new ReadReq } }; + private static DeleteReq DeleteAllRequest(string groupName) => new() + { + Options = new DeleteReq.Types.Options + { + All = new Empty(), + GroupName = groupName + } + }; + private static GetInfoReq GetInfoRequest(string streamName, string groupName) => new() { Options = new GetInfoReq.Types.Options @@ -240,6 +318,15 @@ await call.RequestStream.WriteAsync(new ReadReq } }; + private static GetInfoReq GetAllInfoRequest(string groupName) => new() + { + Options = new GetInfoReq.Types.Options + { + All = new Empty(), + GroupName = groupName + } + }; + private static CreateReq.Types.Settings Settings => new() { CheckpointAfterMs = 10000, diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/PersistentSubscriptionTests/PersistentSubscriptionWithEventNumbersGreaterThan2BillionTests.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/PersistentSubscriptionTests/PersistentSubscriptionWithEventNumbersGreaterThan2BillionTests.cs new file mode 100644 index 0000000000..ddb3fe0bf8 --- /dev/null +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/PersistentSubscriptionTests/PersistentSubscriptionWithEventNumbersGreaterThan2BillionTests.cs @@ -0,0 +1,259 @@ +using System; +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; +using EventStore.Client; +using EventStore.Client.PersistentSubscriptions; +using EventStore.Client.Streams; +using EventStore.Core.Data; +using EventStore.Core.Services.Transport.Grpc; +using EventStore.Core.Tests.Services.Transport.Grpc.StreamsTests; +using Google.Protobuf; +using Grpc.Core; +using NUnit.Framework; +using GrpcMetadata = EventStore.Core.Services.Transport.Grpc.Constants.Metadata; +using ReadReq = EventStore.Client.PersistentSubscriptions.ReadReq; +using ReadResp = EventStore.Client.PersistentSubscriptions.ReadResp; +using Streams = EventStore.Client.Streams.Streams; + +namespace EventStore.Core.Tests.Services.Transport.Grpc.PersistentSubscriptionTests; + +[NonParallelizable] +[TestFixture(typeof(LogFormat.V2), typeof(string))] +public class PersistentSubscriptionWithEventNumbersGreaterThan2BillionTests + : GrpcSpecificationWithExistingRecords +{ + private const long IntMaxValue = int.MaxValue; + private const string StreamName = "persistent-subscription-stream"; + private EventRecord _first; + private EventRecord _second; + private PersistentSubscriptions.PersistentSubscriptionsClient _persistentSubscriptions; + private Streams.StreamsClient _streams; + + public override async ValueTask WriteTestScenario(CancellationToken token) + { + _first = await WriteSingleEvent(StreamName, IntMaxValue + 1, "first", token: token); + _second = await WriteSingleEvent(StreamName, IntMaxValue + 2, "second", token: token); + } + + public override async Task Given() + { + _persistentSubscriptions = new PersistentSubscriptions.PersistentSubscriptionsClient(Channel); + _streams = new Streams.StreamsClient(Channel); + await Append( + "$$" + StreamName, + new[] { ProposedEvent("$metadata", "{\"$tb\":2147483648}") }); + } + + [Test] + public async Task can_create_a_subscription_above_int_max_value() + { + await _persistentSubscriptions.CreateAsync( + CreateRequest(NewName("group"), (ulong)IntMaxValue), + GetCallOptions(AdminCredentials)); + } + + [Test] + public async Task can_update_a_subscription_above_int_max_value() + { + var groupName = NewName("group"); + await _persistentSubscriptions.CreateAsync( + CreateRequest(groupName, 0), GetCallOptions(AdminCredentials)); + + await _persistentSubscriptions.UpdateAsync( + UpdateRequest(groupName, (ulong)IntMaxValue), + GetCallOptions(AdminCredentials)); + } + + [Test] + public async Task delivers_and_appends_events_above_int_max_value() + { + var groupName = NewName("group"); + var thirdId = Guid.NewGuid(); + await _persistentSubscriptions.CreateAsync( + CreateRequest(groupName, (ulong)IntMaxValue), + GetCallOptions(AdminCredentials)); + using var subscription = await Subscribe(groupName); + + await Append( + StreamName, + new[] { ProposedEvent("third", "third", thirdId) }, + (ulong)(IntMaxValue + 2)); + + var expected = new[] + { + (_first.EventId, (ulong)(IntMaxValue + 1)), + (_second.EventId, (ulong)(IntMaxValue + 2)), + (thirdId, (ulong)(IntMaxValue + 3)) + }; + var actual = new List<(Guid eventId, ulong revision)>(); + for (var index = 0; index < expected.Length; index++) + { + Assert.True(await subscription.ResponseStream.MoveNext()); + var response = subscription.ResponseStream.Current; + Assert.AreEqual(ReadResp.ContentOneofCase.Event, response.ContentCase); + actual.Add((Uuid.FromDto(response.Event.Event.Id).ToGuid(), response.Event.Event.StreamRevision)); + await subscription.RequestStream.WriteAsync(new ReadReq + { + Ack = new ReadReq.Types.Ack { Ids = { response.Event.Event.Id } } + }); + } + + CollectionAssert.AreEqual(expected, actual); + } + + [Test] + public async Task resolves_a_link_to_an_event_above_int_max_value() + { + var linkStream = NewName("links"); + var groupName = NewName("group"); + await Append( + linkStream, + new[] { ProposedEvent("$>", $"{IntMaxValue + 1}@{StreamName}") }); + await _persistentSubscriptions.CreateAsync( + CreateRequest(groupName, 0, linkStream, resolveLinks: true), + GetCallOptions(AdminCredentials)); + + using var subscription = await Subscribe(groupName, linkStream); + Assert.True(await subscription.ResponseStream.MoveNext()); + var response = subscription.ResponseStream.Current; + + Assert.AreEqual((ulong)(IntMaxValue + 1), response.Event.Event.StreamRevision); + Assert.AreEqual(_first.EventId, Uuid.FromDto(response.Event.Event.Id).ToGuid()); + Assert.AreEqual(linkStream, response.Event.Link.StreamIdentifier.StreamName.ToStringUtf8()); + } + + private async Task> Subscribe( + string groupName, + string streamName = StreamName) + { + var subscription = _persistentSubscriptions.Read(GetCallOptions(AdminCredentials)); + await subscription.RequestStream.WriteAsync(new ReadReq + { + Options = new ReadReq.Types.Options + { + BufferSize = 1, + GroupName = groupName, + StreamIdentifier = StreamIdentifier(streamName), + UuidOption = new ReadReq.Types.Options.Types.UUIDOption { Structured = new Empty() } + } + }); + Assert.True(await subscription.ResponseStream.MoveNext()); + Assert.AreEqual(ReadResp.ContentOneofCase.SubscriptionConfirmation, + subscription.ResponseStream.Current.ContentCase); + return subscription; + } + + private async Task Append( + string streamName, + IEnumerable events, + ulong? expectedRevision = null) + { + using var call = _streams.BatchAppend(GetCallOptions(AdminCredentials)); + var options = new BatchAppendReq.Types.Options + { + StreamIdentifier = StreamIdentifier(streamName) + }; + if (expectedRevision.HasValue) + { + options.StreamPosition = expectedRevision.Value; + } + else + { + options.Any = new(); + } + + await call.RequestStream.WriteAsync(new BatchAppendReq + { + CorrelationId = Uuid.NewUuid().ToDto(), + IsFinal = true, + Options = options, + ProposedMessages = { events } + }); + await call.RequestStream.CompleteAsync(); + Assert.True(await call.ResponseStream.MoveNext()); + Assert.NotNull(call.ResponseStream.Current.Success); + } + + private static BatchAppendReq.Types.ProposedMessage ProposedEvent( + string type, + string data, + Guid eventId = default) => new() + { + Data = ByteString.CopyFromUtf8(data), + Id = Uuid.FromGuid(eventId == default ? Guid.NewGuid() : eventId).ToDto(), + Metadata = + { + [GrpcMetadata.ContentType] = GrpcMetadata.ContentTypes.ApplicationJson, + [GrpcMetadata.Type] = type + } + }; + + private static CreateReq CreateRequest( + string groupName, + ulong revision, + string streamName = StreamName, + bool resolveLinks = false) => new() + { + Options = new CreateReq.Types.Options + { + GroupName = groupName, + Settings = CreateSettings(resolveLinks), + Stream = new CreateReq.Types.StreamOptions + { + Revision = revision, + StreamIdentifier = StreamIdentifier(streamName) + } + } + }; + + private static UpdateReq UpdateRequest(string groupName, ulong revision) => new() + { + Options = new UpdateReq.Types.Options + { + GroupName = groupName, + Settings = UpdateSettings(), + Stream = new UpdateReq.Types.StreamOptions + { + Revision = revision, + StreamIdentifier = StreamIdentifier(StreamName) + } + } + }; + + private static CreateReq.Types.Settings CreateSettings(bool resolveLinks) => new() + { + CheckpointAfterMs = 100, + ConsumerStrategy = "Pinned", + HistoryBufferSize = 20, + LiveBufferSize = 10, + MaxCheckpointCount = 10, + MaxRetryCount = 10, + MaxSubscriberCount = 1, + MessageTimeoutMs = 10000, + MinCheckpointCount = 1, + ReadBatchSize = 10, + ResolveLinks = resolveLinks + }; + + private static UpdateReq.Types.Settings UpdateSettings() => new() + { + CheckpointAfterMs = 100, + ConsumerStrategy = "Pinned", + HistoryBufferSize = 20, + LiveBufferSize = 10, + MaxCheckpointCount = 10, + MaxRetryCount = 10, + MaxSubscriberCount = 1, + MessageTimeoutMs = 10000, + MinCheckpointCount = 1, + ReadBatchSize = 10 + }; + + private static StreamIdentifier StreamIdentifier(string streamName) => new() + { + StreamName = ByteString.CopyFromUtf8(streamName) + }; + + private static string NewName(string prefix) => $"{prefix}-{Guid.NewGuid():N}"; +} diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/PersistentSubscriptionTests/ReadTests.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/PersistentSubscriptionTests/ReadTests.cs new file mode 100644 index 0000000000..9c882ad386 --- /dev/null +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/PersistentSubscriptionTests/ReadTests.cs @@ -0,0 +1,466 @@ +using System; +using System.Linq; +using System.Threading.Tasks; +using EventStore.Client; +using EventStore.Client.PersistentSubscriptions; +using EventStore.Client.Streams; +using EventStore.Core.Services.Transport.Grpc; +using EventStore.Core.Tests.Services.Transport.Grpc.StreamsTests; +using Google.Protobuf; +using Grpc.Core; +using NUnit.Framework; +using ReadReq = EventStore.Client.PersistentSubscriptions.ReadReq; +using ReadResp = EventStore.Client.PersistentSubscriptions.ReadResp; + +namespace EventStore.Core.Tests.Services.Transport.Grpc.PersistentSubscriptionTests; + +[NonParallelizable] +[TestFixture(typeof(LogFormat.V2), typeof(string))] +public class ReadTests : GrpcSpecification +{ + private PersistentSubscriptions.PersistentSubscriptionsClient _client; + + protected override Task Given() + { + _client = new PersistentSubscriptions.PersistentSubscriptionsClient(Channel); + return Task.CompletedTask; + } + + protected override Task When() => Task.CompletedTask; + + [Test] + public async Task connecting_to_a_missing_subscription_returns_not_found() + { + using var subscription = _client.Read(GetCallOptions(AdminCredentials)); + await subscription.RequestStream.WriteAsync(ReadOptions(NewName("stream"), NewName("group"))); + + var exception = Assert.ThrowsAsync(async () => + await subscription.ResponseStream.MoveNext()); + + Assert.AreEqual(StatusCode.NotFound, exception.StatusCode); + } + + [Test] + public async Task connecting_without_permission_returns_permission_denied() + { + var streamName = $"${NewName("stream")}"; + var groupName = NewName("group"); + await CreateSubscription(streamName, groupName); + + using var subscription = _client.Read(GetCallOptions()); + await subscription.RequestStream.WriteAsync(ReadOptions(streamName, groupName)); + + var exception = Assert.ThrowsAsync(async () => + await subscription.ResponseStream.MoveNext()); + + Assert.AreEqual(StatusCode.PermissionDenied, exception.StatusCode); + } + + [Test] + public async Task connecting_to_a_missing_all_subscription_returns_not_found() + { + using var subscription = _client.Read(GetCallOptions(AdminCredentials)); + await subscription.RequestStream.WriteAsync(ReadAllOptions(NewName("group"))); + + var exception = Assert.ThrowsAsync(async () => + await subscription.ResponseStream.MoveNext()); + + Assert.AreEqual(StatusCode.NotFound, exception.StatusCode); + } + + [Test] + public async Task connecting_to_an_existing_all_subscription_is_confirmed() + { + var groupName = NewName("group"); + await CreateAllSubscription(groupName); + + using var subscription = _client.Read(GetCallOptions(AdminCredentials)); + await subscription.RequestStream.WriteAsync(ReadAllOptions(groupName)); + + Assert.IsTrue(await subscription.ResponseStream.MoveNext()); + Assert.AreEqual( + ReadResp.ContentOneofCase.SubscriptionConfirmation, + subscription.ResponseStream.Current.ContentCase); + } + + [Test] + public async Task connecting_to_an_all_subscription_without_permission_returns_permission_denied() + { + var groupName = NewName("group"); + await CreateAllSubscription(groupName); + + using var subscription = _client.Read(GetCallOptions()); + await subscription.RequestStream.WriteAsync(ReadAllOptions(groupName)); + + var exception = Assert.ThrowsAsync(async () => + await subscription.ResponseStream.MoveNext()); + + Assert.AreEqual(StatusCode.PermissionDenied, exception.StatusCode); + } + + [Test] + public async Task connecting_beyond_the_subscriber_limit_returns_failed_precondition() + { + var streamName = NewName("stream"); + var groupName = NewName("group"); + await CreateSubscription(streamName, groupName, maxSubscriberCount: 1); + + using var first = await Subscribe(streamName, groupName); + using var second = _client.Read(GetCallOptions(AdminCredentials)); + await second.RequestStream.WriteAsync(ReadOptions(streamName, groupName)); + + var exception = Assert.ThrowsAsync(async () => + await second.ResponseStream.MoveNext()); + + Assert.AreEqual(StatusCode.FailedPrecondition, exception.StatusCode); + } + + [Test] + public async Task subscription_starts_at_the_requested_revision_and_acknowledges_the_event() + { + var streamName = NewName("stream"); + var groupName = NewName("group"); + await Append(streamName, CreateEvents(3).ToArray()); + await CreateSubscription(streamName, groupName, revision: 1); + + using var subscription = await Subscribe(streamName, groupName, bufferSize: 1); + var response = await ReadEvent(subscription); + + Assert.AreEqual(1, response.Event.Event.StreamRevision); + + await subscription.RequestStream.WriteAsync(new ReadReq + { + Ack = new ReadReq.Types.Ack { Ids = { response.Event.Event.Id } } + }); + + var next = await ReadEvent(subscription); + Assert.AreEqual(2, next.Event.Event.StreamRevision); + } + + [Test] + public async Task subscription_from_end_receives_only_new_events() + { + var streamName = NewName("stream"); + var groupName = NewName("group"); + await Append(streamName, CreateEvent()); + await CreateSubscription(streamName, groupName, startFromEnd: true); + + using var subscription = await Subscribe(streamName, groupName); + await Append(streamName, CreateEvent()); + var response = await ReadEvent(subscription); + + Assert.AreEqual(1, response.Event.Event.StreamRevision); + } + + [Test] + public async Task subscription_on_a_missing_stream_receives_the_first_live_event() + { + var streamName = NewName("stream"); + var groupName = NewName("group"); + var proposedEvent = CreateEvent(); + await CreateSubscription(streamName, groupName); + + using var subscription = await Subscribe(streamName, groupName); + await Append(streamName, proposedEvent); + var response = await ReadEvent(subscription); + + Assert.AreEqual(0, response.Event.Event.StreamRevision); + Assert.AreEqual(Uuid.FromDto(proposedEvent.Id), Uuid.FromDto(response.Event.Event.Id)); + } + + [Test] + public async Task subscription_waits_for_a_requested_future_revision() + { + var streamName = NewName("stream"); + var groupName = NewName("group"); + await Append(streamName, CreateEvents(3).ToArray()); + await CreateSubscription(streamName, groupName, revision: 3); + + using var subscription = await Subscribe(streamName, groupName); + await Append(streamName, CreateEvent()); + var response = await ReadEvent(subscription); + + Assert.AreEqual(3, response.Event.Event.StreamRevision); + } + + [Test] + public async Task manual_acknowledgement_drains_multiple_buffer_windows() + { + const int eventCount = 20; + var streamName = NewName("stream"); + var groupName = NewName("group"); + await CreateSubscription(streamName, groupName, startFromEnd: true); + + using var subscription = await Subscribe(streamName, groupName, bufferSize: 5); + await Append(streamName, CreateEvents(eventCount).ToArray()); + for (var revision = 0; revision < eventCount; revision++) + { + var response = await ReadEvent(subscription); + Assert.AreEqual((ulong)revision, response.Event.Event.StreamRevision); + await subscription.RequestStream.WriteAsync(new ReadReq + { + Ack = new ReadReq.Types.Ack { Ids = { response.Event.Event.Id } } + }); + } + } + + [Test] + public async Task retrying_a_nacked_event_preserves_each_retry_count_until_acknowledged() + { + var streamName = NewName("stream"); + var groupName = NewName("group"); + await Append(streamName, CreateEvent()); + await CreateSubscription(streamName, groupName); + + using var subscription = await Subscribe(streamName, groupName); + var response = await ReadEvent(subscription); + var eventId = response.Event.Event.Id; + for (var retryCount = 1; retryCount <= 5; retryCount++) + { + await subscription.RequestStream.WriteAsync(new ReadReq + { + Nack = new ReadReq.Types.Nack + { + Action = ReadReq.Types.Nack.Types.Action.Retry, + Ids = { eventId }, + Reason = "retry" + } + }); + response = await ReadEvent(subscription); + + Assert.AreEqual(eventId, response.Event.Event.Id); + Assert.AreEqual(retryCount, response.Event.RetryCount); + } + + await subscription.RequestStream.WriteAsync(new ReadReq + { + Ack = new ReadReq.Types.Ack { Ids = { eventId } } + }); + } + + [Test] + public async Task disconnected_event_with_no_retries_is_not_redelivered() + { + var streamName = NewName("stream"); + var groupName = NewName("group"); + await CreateSubscription(streamName, groupName, startFromEnd: true, maxRetryCount: 0); + var firstEvent = CreateEvent(); + var firstSubscription = await Subscribe(streamName, groupName, bufferSize: 1); + await Append(streamName, firstEvent); + await ReadEvent(firstSubscription); + firstSubscription.Dispose(); + + await WaitForNoSubscribers(streamName, groupName); + + using var replacement = await Subscribe(streamName, groupName, bufferSize: 1); + var nextEvent = CreateEvent(); + await Append(streamName, nextEvent); + var response = await ReadEvent(replacement); + + Assert.AreEqual(Uuid.FromDto(nextEvent.Id), Uuid.FromDto(response.Event.Event.Id)); + } + + [Test] + public async Task link_resolution_supports_stream_names_containing_at_symbols() + { + var targetStream = $"target@{NewName("domain")}"; + var linkStream = NewName("links"); + var groupName = NewName("group"); + var target = CreateEvent("target-event"); + target.Data = ByteString.CopyFromUtf8("data"); + await Append(targetStream, target); + await Append(linkStream, LinkTo(0, targetStream)); + await CreateSubscription(linkStream, groupName, resolveLinks: true); + + using var subscription = await Subscribe(linkStream, groupName); + var response = await ReadEvent(subscription); + + Assert.AreEqual(targetStream, response.Event.Event.StreamIdentifier.StreamName.ToStringUtf8()); + Assert.AreEqual("data", response.Event.Event.Data.ToStringUtf8()); + Assert.AreEqual(linkStream, response.Event.Link.StreamIdentifier.StreamName.ToStringUtf8()); + } + + private async Task CreateSubscription( + string streamName, + string groupName, + ulong? revision = null, + bool startFromEnd = false, + bool resolveLinks = false, + int maxSubscriberCount = 40, + int maxRetryCount = 10) + { + var streamOptions = new CreateReq.Types.StreamOptions + { + StreamIdentifier = StreamIdentifier(streamName) + }; + if (revision.HasValue) + { + streamOptions.Revision = revision.Value; + } + else if (startFromEnd) + { + streamOptions.End = new Empty(); + } + else + { + streamOptions.Start = new Empty(); + } + + await _client.CreateAsync(new CreateReq + { + Options = new CreateReq.Types.Options + { + GroupName = groupName, + Stream = streamOptions, + Settings = new CreateReq.Types.Settings + { + CheckpointAfterMs = 100, + ConsumerStrategy = "Pinned", + HistoryBufferSize = 20, + LiveBufferSize = 10, + MaxCheckpointCount = 10, + MaxRetryCount = maxRetryCount, + MaxSubscriberCount = maxSubscriberCount, + MessageTimeoutMs = 10000, + MinCheckpointCount = 1, + ReadBatchSize = 10, + ResolveLinks = resolveLinks + } + } + }, GetCallOptions(AdminCredentials)); + } + + private async Task CreateAllSubscription(string groupName) + { + await _client.CreateAsync(new CreateReq + { + Options = new CreateReq.Types.Options + { + GroupName = groupName, + All = new CreateReq.Types.AllOptions + { + Start = new Empty(), + NoFilter = new Empty() + }, + Settings = new CreateReq.Types.Settings + { + CheckpointAfterMs = 100, + ConsumerStrategy = "Pinned", + HistoryBufferSize = 20, + LiveBufferSize = 10, + MaxCheckpointCount = 10, + MaxRetryCount = 10, + MaxSubscriberCount = 40, + MessageTimeoutMs = 10000, + MinCheckpointCount = 1, + ReadBatchSize = 10 + } + } + }, GetCallOptions(AdminCredentials)); + } + + private async Task> Subscribe( + string streamName, + string groupName, + int bufferSize = 10) + { + var subscription = _client.Read(GetCallOptions(AdminCredentials)); + await subscription.RequestStream.WriteAsync(ReadOptions(streamName, groupName, bufferSize)); + if (!await subscription.ResponseStream.MoveNext() || + subscription.ResponseStream.Current.ContentCase != ReadResp.ContentOneofCase.SubscriptionConfirmation) + { + subscription.Dispose(); + throw new InvalidOperationException("Persistent subscription was not confirmed."); + } + + return subscription; + } + + private static async Task ReadEvent(AsyncDuplexStreamingCall subscription) + { + if (!await subscription.ResponseStream.MoveNext() || + subscription.ResponseStream.Current.ContentCase != ReadResp.ContentOneofCase.Event) + { + throw new InvalidOperationException("Persistent subscription did not return an event."); + } + + return subscription.ResponseStream.Current; + } + + private async Task WaitForNoSubscribers(string streamName, string groupName) + { + var deadline = DateTime.UtcNow.AddSeconds(10); + do + { + var response = await _client.GetInfoAsync(new GetInfoReq + { + Options = new GetInfoReq.Types.Options + { + GroupName = groupName, + StreamIdentifier = StreamIdentifier(streamName) + } + }, GetCallOptions(AdminCredentials)); + if (response.SubscriptionInfo.Connections.Count == 0) + { + return; + } + + await Task.Delay(10); + } while (DateTime.UtcNow < deadline); + + Assert.Fail("Disconnected subscriber remained registered."); + } + + private async Task Append(string streamName, params BatchAppendReq.Types.ProposedMessage[] events) + { + var response = await AppendToStreamBatch(new BatchAppendReq + { + CorrelationId = Uuid.NewUuid().ToDto(), + IsFinal = true, + Options = new BatchAppendReq.Types.Options + { + Any = new(), + StreamIdentifier = StreamIdentifier(streamName) + }, + ProposedMessages = { events } + }); + + Assert.NotNull(response.Success); + } + + private static BatchAppendReq.Types.ProposedMessage LinkTo(ulong revision, string streamName) + { + var link = CreateEvent("$>"); + link.Data = ByteString.CopyFromUtf8($"{revision}@{streamName}"); + return link; + } + + private static ReadReq ReadOptions(string streamName, string groupName, int bufferSize = 10) => new() + { + Options = new ReadReq.Types.Options + { + BufferSize = bufferSize, + GroupName = groupName, + StreamIdentifier = StreamIdentifier(streamName), + UuidOption = new ReadReq.Types.Options.Types.UUIDOption { Structured = new Empty() } + } + }; + + private static ReadReq ReadAllOptions(string groupName, int bufferSize = 10) => new() + { + Options = new ReadReq.Types.Options + { + All = new Empty(), + BufferSize = bufferSize, + GroupName = groupName, + UuidOption = new ReadReq.Types.Options.Types.UUIDOption { Structured = new Empty() } + } + }; + + private static StreamIdentifier StreamIdentifier(string streamName) => new() + { + StreamName = ByteString.CopyFromUtf8(streamName) + }; + + private static string NewName(string prefix) => $"{prefix}-{Guid.NewGuid():N}"; +} diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/GrpcSpecificationWithExistingRecords.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/GrpcSpecificationWithExistingRecords.cs new file mode 100644 index 0000000000..83839d9e61 --- /dev/null +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/GrpcSpecificationWithExistingRecords.cs @@ -0,0 +1,157 @@ +using System; +using System.IO; +using System.Text; +using System.Threading; +using System.Threading.Tasks; +using EventStore.Common.Utils; +using EventStore.Core.Data; +using EventStore.Core.Helpers; +using EventStore.Core.LogAbstraction; +using EventStore.Core.Tests.Helpers; +using EventStore.Core.Tests.TransactionLog; +using EventStore.Core.TransactionLog.Checkpoint; +using EventStore.Core.TransactionLog.Chunks; +using EventStore.Core.TransactionLog.LogRecords; +using Grpc.Core; +using Grpc.Net.Client; +using NUnit.Framework; + +namespace EventStore.Core.Tests.Services.Transport.Grpc.StreamsTests; + +public abstract class GrpcSpecificationWithExistingRecords + : SpecificationWithDirectoryPerTestFixture +{ + private string _dbPath; + private LogFormatAbstractor _logFormatFactory; + private TFChunkDb _db; + private TFChunkWriter _writer; + private ICheckpoint _writerCheckpoint; + private ICheckpoint _chaserCheckpoint; + + protected MiniNode Node; + protected GrpcChannel Channel; + + protected static (string userName, string password) AdminCredentials => ("admin", "changeit"); + + public abstract ValueTask WriteTestScenario(CancellationToken token); + + public abstract Task Given(); + + [OneTimeSetUp] + public override async Task TestFixtureSetUp() + { + await base.TestFixtureSetUp(); + _dbPath = Path.Combine(PathName, $"mini-node-db-{Guid.NewGuid():N}"); + _logFormatFactory = LogFormatHelper.LogFormatFactory.Create(new() + { + IndexDirectory = GetFilePathFor("index") + }); + + Directory.CreateDirectory(_dbPath); + + _writerCheckpoint = new MemoryMappedFileCheckpoint( + Path.Combine(_dbPath, Checkpoint.Writer + ".chk"), Checkpoint.Writer); + _chaserCheckpoint = new MemoryMappedFileCheckpoint( + Path.Combine(_dbPath, Checkpoint.Chaser + ".chk"), Checkpoint.Chaser); + _db = new TFChunkDb(TFChunkHelper.CreateDbConfig( + _dbPath, _writerCheckpoint, _chaserCheckpoint, TFConsts.ChunkSize)); + await _db.Open(); + + _writer = new TFChunkWriter(_db); + _writer.Open(); + var partitionManager = _logFormatFactory.CreatePartitionManager( + new TFChunkReader(_db, _writerCheckpoint), _writer); + await partitionManager.Initialize(CancellationToken.None); + await WriteTestScenario(CancellationToken.None); + + await _writer.DisposeAsync(); + _writer = null; + _writerCheckpoint.Flush(); + _chaserCheckpoint.Write(_writerCheckpoint.Read()); + _chaserCheckpoint.Flush(); + await _db.DisposeAsync(); + _db = null; + + Node = new MiniNode(PathName, dbPath: _dbPath); + await Node.Start(); + await Node.AdminUserCreated; + Channel = GrpcChannel.ForAddress(new UriBuilder { Scheme = Uri.UriSchemeHttps }.Uri, + new GrpcChannelOptions + { + HttpClient = Node.HttpClient, + DisposeHttpClient = false + }); + + await Given().WithTimeout(TimeSpan.FromSeconds(30)); + } + + [OneTimeTearDown] + public override async Task TestFixtureTearDown() + { + Channel?.Dispose(); + _logFormatFactory?.Dispose(); + if (Node is not null) + { + await Node.Shutdown(); + } + await base.TestFixtureTearDown(); + } + + protected CallOptions GetCallOptions((string userName, string password) credentials) => + new(credentials: CallCredentials.FromInterceptor((_, metadata) => + { + var value = Convert.ToBase64String( + Encoding.ASCII.GetBytes($"{credentials.userName}:{credentials.password}")); + metadata.Add(new Metadata.Entry("authorization", $"Basic {value}")); + return Task.CompletedTask; + })); + + protected async ValueTask WriteSingleEvent( + string eventStreamName, + long eventNumber, + string data, + Guid eventId = default, + string eventType = "some-type", + CancellationToken token = default) + { + var position = _writer.Position; + _logFormatFactory.StreamNameIndex.GetOrReserve( + _logFormatFactory.RecordFactory, + eventStreamName, + position, + out var eventStreamId, + out var streamRecord); + if (streamRecord is not null) + { + (_, position) = await _writer.Write(streamRecord, token); + } + + _logFormatFactory.EventTypeIndex.GetOrReserveEventType( + _logFormatFactory.RecordFactory, + eventType, + position, + out var eventTypeId, + out var eventTypeRecord); + if (eventTypeRecord is not null) + { + (_, position) = await _writer.Write(eventTypeRecord, token); + } + + var prepare = LogRecord.SingleWrite( + _logFormatFactory.RecordFactory, + position, + eventId == default ? Guid.NewGuid() : eventId, + Guid.NewGuid(), + eventStreamId, + eventNumber - 1, + eventTypeId, + Helper.UTF8NoBom.GetBytes(data), + null); + var (written, nextPosition) = await _writer.Write(prepare, token); + Assert.IsTrue(written); + var commit = LogRecord.Commit(nextPosition, prepare.CorrelationId, prepare.LogPosition, eventNumber); + Assert.IsTrue(await _writer.Write(commit, token) is (true, _)); + + return new EventRecord(eventNumber, prepare, eventStreamName, eventType); + } +} diff --git a/src/EventStore.Core/Messages/ClientMessage.cs b/src/EventStore.Core/Messages/ClientMessage.cs index 905652eca3..3b0f349bfa 100644 --- a/src/EventStore.Core/Messages/ClientMessage.cs +++ b/src/EventStore.Core/Messages/ClientMessage.cs @@ -1928,10 +1928,20 @@ public CheckpointReached(Guid correlationId, TFPos? position) [DerivedMessage(CoreMessage.Client)] public partial class UnsubscribeFromStream : ReadRequestMessage { + public enum SubscriptionEndReason + { + Unsubscribed, + ConnectionClosed + } + + public readonly SubscriptionEndReason Reason; + public UnsubscribeFromStream(Guid internalCorrId, Guid correlationId, IEnvelope envelope, - ClaimsPrincipal user, DateTime? expires = null) + ClaimsPrincipal user, SubscriptionEndReason reason = SubscriptionEndReason.Unsubscribed, + DateTime? expires = null) : base(internalCorrId, correlationId, envelope, user, expires) { + Reason = reason; } } diff --git a/src/EventStore.Core/Services/PersistentSubscription/PersistentSubscriptionService.cs b/src/EventStore.Core/Services/PersistentSubscription/PersistentSubscriptionService.cs index a1c81c8374..c1d3213f9d 100644 --- a/src/EventStore.Core/Services/PersistentSubscription/PersistentSubscriptionService.cs +++ b/src/EventStore.Core/Services/PersistentSubscription/PersistentSubscriptionService.cs @@ -218,6 +218,12 @@ public void Handle(ClientMessage.UnsubscribeFromStream message) return; } + if (message.Reason == ClientMessage.UnsubscribeFromStream.SubscriptionEndReason.ConnectionClosed) + { + DisconnectFromStream(message.CorrelationId); + return; + } + UnsubscribeFromStream(message.CorrelationId, true); } @@ -1121,6 +1127,18 @@ public void Handle(TcpMessage.ConnectionClosed message) } } + private void DisconnectFromStream(Guid connectionId) + { + foreach (var subscription in _subscriptionsById.Values) + { + if (subscription.RemoveClientByConnectionId(connectionId)) + { + Log.Debug("Persistent subscription {subscription} lost connection {connectionId}", + subscription.SubscriptionId, connectionId); + } + } + } + public async ValueTask ConnectToPersistentSubscription( IPersistentSubscriptionEventSource eventSource, string groupName, diff --git a/src/EventStore.Core/Services/Transport/Grpc/PersistentSubscriptions.Read.cs b/src/EventStore.Core/Services/Transport/Grpc/PersistentSubscriptions.Read.cs index cb108471ea..d374a3fcf3 100644 --- a/src/EventStore.Core/Services/Transport/Grpc/PersistentSubscriptions.Read.cs +++ b/src/EventStore.Core/Services/Transport/Grpc/PersistentSubscriptions.Read.cs @@ -339,7 +339,8 @@ void Fail(Exception ex) public ValueTask DisposeAsync() { _publisher.Publish(new ClientMessage.UnsubscribeFromStream(Guid.NewGuid(), _correlationId, - new NoopEnvelope(), _user)); + new NoopEnvelope(), _user, + ClientMessage.UnsubscribeFromStream.SubscriptionEndReason.ConnectionClosed)); _channel.Writer.TryComplete(); return new ValueTask(Task.CompletedTask); }