diff --git a/src/Squidex.Infrastructure/CQRS/Events/Actors/EventConsumerActor.cs b/src/Squidex.Infrastructure/CQRS/Events/Actors/EventConsumerActor.cs index 210435b27..2c2f3c3ed 100644 --- a/src/Squidex.Infrastructure/CQRS/Events/Actors/EventConsumerActor.cs +++ b/src/Squidex.Infrastructure/CQRS/Events/Actors/EventConsumerActor.cs @@ -31,11 +31,6 @@ namespace Squidex.Infrastructure.CQRS.Events.Actors public IEventConsumer EventConsumer { get; set; } } - private sealed class StopFailed - { - public Exception Exception { get; set; } - } - private abstract class SubscriptionMessage { public IEventSubscription Subscription { get; set; } @@ -84,16 +79,16 @@ namespace Squidex.Infrastructure.CQRS.Events.Actors } } - protected override Task OnError(Exception exception) + protected override async Task OnError(Exception exception) { log.LogError(exception, w => w .WriteProperty("action", "HandleEvent") .WriteProperty("state", "Failed") .WriteProperty("eventConsumer", eventConsumer.Name)); - DispatchAsync(new StopFailed { Exception = exception }).Forget(); + await StopAsync(exception); - return TaskHelper.Done; + isRunning = false; } Task IEventSubscriber.OnEventAsync(IEventSubscription subscription, StoredEvent @event) @@ -135,15 +130,6 @@ namespace Squidex.Infrastructure.CQRS.Events.Actors break; } - case StopFailed stopFailed when isSetup && isRunning: - { - await StopAsync(stopFailed.Exception); - - isRunning = false; - - break; - } - case StopConsumerMessage stopConsumer when isSetup && isRunning: { await StopAsync(); diff --git a/tests/Squidex.Infrastructure.Tests/CQRS/Events/Actors/EventConsumerActorTests.cs b/tests/Squidex.Infrastructure.Tests/CQRS/Events/Actors/EventConsumerActorTests.cs index 23550364c..376e63ad7 100644 --- a/tests/Squidex.Infrastructure.Tests/CQRS/Events/Actors/EventConsumerActorTests.cs +++ b/tests/Squidex.Infrastructure.Tests/CQRS/Events/Actors/EventConsumerActorTests.cs @@ -196,9 +196,34 @@ namespace Squidex.Infrastructure.CQRS.Events.Actors var @event = new StoredEvent(Guid.NewGuid().ToString(), 123, eventData); await SubscribeAsync(); + await sutSubscriber.OnEventAsync(eventSubscription, @event); + + sut.Dispose(); + + A.CallTo(() => eventConsumer.On(envelope)) + .MustHaveHappened(); + + A.CallTo(() => eventConsumerInfoRepository.SetPositionAsync(consumerName, @event.EventPosition, false)) + .MustNotHaveHappened(); + A.CallTo(() => eventConsumerInfoRepository.StopAsync(consumerName, exception.ToString())) + .MustHaveHappened(); + } + + [Fact] + public async Task Should_start_after_stop_when_handling_failed() + { + var exception = new InvalidOperationException("Exception"); + + A.CallTo(() => eventConsumer.On(envelope)) + .Throws(exception); + + var @event = new StoredEvent(Guid.NewGuid().ToString(), 123, eventData); + + await SubscribeAsync(); await sutSubscriber.OnEventAsync(eventSubscription, @event); + sutActor.Tell(new StartConsumerMessage()); sut.Dispose(); A.CallTo(() => eventConsumer.On(envelope)) @@ -209,6 +234,9 @@ namespace Squidex.Infrastructure.CQRS.Events.Actors A.CallTo(() => eventConsumerInfoRepository.StopAsync(consumerName, exception.ToString())) .MustHaveHappened(); + + A.CallTo(() => eventConsumerInfoRepository.StartAsync(consumerName)) + .MustHaveHappened(); } [Fact] @@ -222,7 +250,6 @@ namespace Squidex.Infrastructure.CQRS.Events.Actors var @event = new StoredEvent(Guid.NewGuid().ToString(), 123, eventData); await SubscribeAsync(); - await sutSubscriber.OnEventAsync(eventSubscription, @event); sut.Dispose();