diff --git a/src/Squidex.Infrastructure/CQRS/Events/Actors/EventConsumerActor.cs b/src/Squidex.Infrastructure/CQRS/Events/Actors/EventConsumerActor.cs index eb98f5824..b28b6d2e6 100644 --- a/src/Squidex.Infrastructure/CQRS/Events/Actors/EventConsumerActor.cs +++ b/src/Squidex.Infrastructure/CQRS/Events/Actors/EventConsumerActor.cs @@ -100,6 +100,14 @@ namespace Squidex.Infrastructure.CQRS.Events.Actors } } + public async Task WaitForCompletionAsync() + { + while (dispatcher.InputCount > 0) + { + await Task.Delay(20); + } + } + public Task SubscribeAsync(IEventConsumer eventConsumer) { Guard.NotNull(eventConsumer, nameof(eventConsumer)); diff --git a/tests/Squidex.Infrastructure.Tests/CQRS/Commands/AggregateHandlerTests.cs b/tests/Squidex.Infrastructure.Tests/CQRS/Commands/AggregateHandlerTests.cs index c2cf61f7b..fd388e909 100644 --- a/tests/Squidex.Infrastructure.Tests/CQRS/Commands/AggregateHandlerTests.cs +++ b/tests/Squidex.Infrastructure.Tests/CQRS/Commands/AggregateHandlerTests.cs @@ -87,7 +87,7 @@ namespace Squidex.Infrastructure.CQRS.Commands await sut.CreateAsync(context, async x => { - await Task.Delay(1); + await Task.Yield(); passedDomainObject = x; }); @@ -139,7 +139,7 @@ namespace Squidex.Infrastructure.CQRS.Commands await sut.UpdateAsync(context, async x => { - await Task.Delay(1); + await Task.Yield(); passedDomainObject = x; }); diff --git a/tests/Squidex.Infrastructure.Tests/CQRS/Events/Actors/EventConsumerActorTests.cs b/tests/Squidex.Infrastructure.Tests/CQRS/Events/Actors/EventConsumerActorTests.cs index f9bc96132..dd9b71c10 100644 --- a/tests/Squidex.Infrastructure.Tests/CQRS/Events/Actors/EventConsumerActorTests.cs +++ b/tests/Squidex.Infrastructure.Tests/CQRS/Events/Actors/EventConsumerActorTests.cs @@ -207,6 +207,10 @@ namespace Squidex.Infrastructure.CQRS.Events.Actors await OnSubscribeAsync(); await OnErrorAsync(eventSubscription, ex); + await Task.Delay(200); + + await sut.WaitForCompletionAsync(); + sut.Dispose(); A.CallTo(() => eventConsumerInfoRepository.SetAsync(consumerName, consumerInfo.Position, false, null)) @@ -341,25 +345,19 @@ namespace Squidex.Infrastructure.CQRS.Events.Actors .MustHaveHappened(Repeated.Exactly.Twice); } - private async Task OnErrorAsync(IEventSubscription subscriber, Exception ex) + private Task OnErrorAsync(IEventSubscription subscriber, Exception ex) { - await sutSubscriber.OnErrorAsync(subscriber, ex); - - await Task.Delay(200); + return sutSubscriber.OnErrorAsync(subscriber, ex); } - private async Task OnEventAsync(IEventSubscription subscriber, StoredEvent ev) + private Task OnEventAsync(IEventSubscription subscriber, StoredEvent ev) { - await sutSubscriber.OnEventAsync(subscriber, ev); - - await Task.Delay(200); + return sutSubscriber.OnEventAsync(subscriber, ev); } - private async Task OnSubscribeAsync() + private Task OnSubscribeAsync() { - await sut.SubscribeAsync(eventConsumer); - - await Task.Delay(200); + return sut.SubscribeAsync(eventConsumer); } } } \ No newline at end of file diff --git a/tests/Squidex.Infrastructure.Tests/CQRS/Events/PollingSubscriptionTests.cs b/tests/Squidex.Infrastructure.Tests/CQRS/Events/PollingSubscriptionTests.cs new file mode 100644 index 000000000..21c74e714 --- /dev/null +++ b/tests/Squidex.Infrastructure.Tests/CQRS/Events/PollingSubscriptionTests.cs @@ -0,0 +1,68 @@ +// ========================================================================== +// PollingSubscriptionTests.cs +// Squidex Headless CMS +// ========================================================================== +// Copyright (c) Squidex Group +// All rights reserved. +// ========================================================================== + +using System; +using System.Threading; +using System.Threading.Tasks; +using FakeItEasy; +using Xunit; + +namespace Squidex.Infrastructure.CQRS.Events +{ + public class PollingSubscriptionTests + { + private readonly IEventStore eventStore = A.Fake(); + private readonly IEventNotifier eventNotifier = new DefaultEventNotifier(new InMemoryPubSub()); + private readonly IEventSubscriber eventSubscriber = A.Fake(); + private readonly PollingSubscription sut; + private readonly string position = Guid.NewGuid().ToString(); + + public PollingSubscriptionTests() + { + sut = new PollingSubscription(eventStore, eventNotifier, eventSubscriber, "^my-stream", position); + } + + [Fact] + public async Task Should_subscribe_on_start() + { + await WaitAndStopAsync(); + + A.CallTo(() => eventStore.GetEventsAsync(A>.Ignored, A.Ignored, "^my-stream", position)) + .MustHaveHappened(Repeated.Exactly.Once); + } + + [Fact] + public async Task Should_not_subscribe_on_notify_when_stream_matches() + { + eventNotifier.NotifyEventsStored("other-stream-123"); + + await WaitAndStopAsync(); + + A.CallTo(() => eventStore.GetEventsAsync(A>.Ignored, A.Ignored, "^my-stream", position)) + .MustHaveHappened(Repeated.Exactly.Once); + } + + [Fact] + public async Task Should_subscribe_on_notify_when_stream_matches() + { + eventNotifier.NotifyEventsStored("my-stream-123"); + + await WaitAndStopAsync(); + + A.CallTo(() => eventStore.GetEventsAsync(A>.Ignored, A.Ignored, "^my-stream", position)) + .MustHaveHappened(Repeated.Exactly.Twice); + } + + private async Task WaitAndStopAsync() + { + await Task.Delay(1000); + + await sut.StopAsync(); + } + } +} diff --git a/tests/Squidex.Infrastructure.Tests/UsageTracking/BackgroundUsageTrackerTests.cs b/tests/Squidex.Infrastructure.Tests/UsageTracking/BackgroundUsageTrackerTests.cs index 1e2bef9d5..6259fd841 100644 --- a/tests/Squidex.Infrastructure.Tests/UsageTracking/BackgroundUsageTrackerTests.cs +++ b/tests/Squidex.Infrastructure.Tests/UsageTracking/BackgroundUsageTrackerTests.cs @@ -112,8 +112,7 @@ namespace Squidex.Infrastructure.UsageTracking await sut.TrackAsync("key1", 0, 1000); sut.Next(); - - await Task.Delay(100); + sut.Dispose(); A.CallTo(() => usageStore.TrackUsagesAsync(A.Ignored, A.Ignored, A.Ignored, A.Ignored)).MustNotHaveHappened(); }