diff --git a/backend/src/Squidex.Domain.Apps.Entities/Assets/AssetUsageTracker_EventHandling.cs b/backend/src/Squidex.Domain.Apps.Entities/Assets/AssetUsageTracker_EventHandling.cs index d7270290b..b4e32f3dc 100644 --- a/backend/src/Squidex.Domain.Apps.Entities/Assets/AssetUsageTracker_EventHandling.cs +++ b/backend/src/Squidex.Domain.Apps.Entities/Assets/AssetUsageTracker_EventHandling.cs @@ -5,14 +5,12 @@ // All rights reserved. Licensed under the MIT license. // ========================================================================== -using System.Text.Json; using Microsoft.Extensions.Caching.Memory; using Microsoft.Extensions.Options; using Squidex.Domain.Apps.Core.Tags; using Squidex.Domain.Apps.Events.Assets; using Squidex.Infrastructure; using Squidex.Infrastructure.EventSourcing; -using Squidex.Infrastructure.Json.System; using Squidex.Infrastructure.States; using Squidex.Infrastructure.UsageTracking; @@ -151,11 +149,6 @@ namespace Squidex.Domain.Apps.Entities.Assets await tagService.UpdateAsync(appId, TagGroups.Assets, updates); } - var options = new JsonSerializerOptions(); - options.Converters.Add(new StringConverter()); - - Console.WriteLine($"Writing tags: {JsonSerializer.Serialize(tagsPerApp, options)}"); - await store.WriteManyAsync(tagsPerAsset.Select(x => new SnapshotWriteJob(x.Key, x.Value, 0))); } diff --git a/backend/src/Squidex.Infrastructure/EventSourcing/PollingSubscription.cs b/backend/src/Squidex.Infrastructure/EventSourcing/PollingSubscription.cs index 927240c74..4f732dd24 100644 --- a/backend/src/Squidex.Infrastructure/EventSourcing/PollingSubscription.cs +++ b/backend/src/Squidex.Infrastructure/EventSourcing/PollingSubscription.cs @@ -12,8 +12,48 @@ namespace Squidex.Infrastructure.EventSourcing { public sealed class PollingSubscription : IEventSubscription { + private readonly RecentEvents recentEvents = new RecentEvents(); private readonly CompletionTimer timer; + private sealed class RecentEvents + { + private const int Capacity = 50; + private readonly HashSet eventIds = new HashSet(Capacity); + private readonly Queue<(Guid, string)> eventQueue = new Queue<(Guid, string)>(Capacity); + + public string? FirstPosition() + { + if (eventQueue.Count == 0) + { + return null; + } + + return eventQueue.Peek().Item2; + } + + public bool Add(StoredEvent @event) + { + var id = @event.Data.Headers.EventId(); + + if (eventIds.Contains(id)) + { + return false; + } + + while (eventQueue.Count >= Capacity) + { + var (storedId, _) = eventQueue.Dequeue(); + + eventIds.Remove(storedId); + } + + eventIds.Add(id); + eventQueue.Enqueue((id, @event.EventPosition)); + + return true; + } + } + public PollingSubscription( IEventStore eventStore, IEventSubscriber eventSubscriber, @@ -24,12 +64,23 @@ namespace Squidex.Infrastructure.EventSourcing { try { - await foreach (var storedEvent in eventStore.QueryAllAsync(streamFilter, position, ct: ct)) + var newEventCount = 0; + do { - await eventSubscriber.OnNextAsync(this, storedEvent); + newEventCount = 0; + + await foreach (var storedEvent in eventStore.QueryAllAsync(streamFilter, position, ct: ct)) + { + if (recentEvents.Add(storedEvent)) + { + await eventSubscriber.OnNextAsync(this, storedEvent); + newEventCount++; + } - position = storedEvent.EventPosition; + position = recentEvents.FirstPosition(); + } } + while (newEventCount > 50); } catch (Exception ex) { diff --git a/backend/tests/Squidex.Infrastructure.Tests/EventSourcing/PollingSubscriptionTests.cs b/backend/tests/Squidex.Infrastructure.Tests/EventSourcing/PollingSubscriptionTests.cs index e68fbb0ce..a6a75ce44 100644 --- a/backend/tests/Squidex.Infrastructure.Tests/EventSourcing/PollingSubscriptionTests.cs +++ b/backend/tests/Squidex.Infrastructure.Tests/EventSourcing/PollingSubscriptionTests.cs @@ -5,7 +5,9 @@ // All rights reserved. Licensed under the MIT license. // ========================================================================== +using System.Globalization; using FakeItEasy; +using FluentAssertions; using Xunit; namespace Squidex.Infrastructure.EventSourcing @@ -28,7 +30,7 @@ namespace Squidex.Infrastructure.EventSourcing } [Fact] - public async Task Should_propagate_exception_to_subscriber() + public async Task Should_forward_exception_to_subscriber() { var ex = new InvalidOperationException(); @@ -44,7 +46,7 @@ namespace Squidex.Infrastructure.EventSourcing } [Fact] - public async Task Should_propagate_operation_cancelled_exception_to_subscriber() + public async Task Should_forward_operation_cancelled_exception_to_subscriber() { var ex = new OperationCanceledException(); @@ -60,7 +62,7 @@ namespace Squidex.Infrastructure.EventSourcing } [Fact] - public async Task Should_propagate_aggregate_operation_cancelled_exception_to_subscriber() + public async Task Should_forward_aggregate_operation_cancelled_exception_to_subscriber() { var ex = new AggregateException(new OperationCanceledException()); @@ -88,6 +90,69 @@ namespace Squidex.Infrastructure.EventSourcing .MustHaveHappened(2, Times.Exactly); } + [Fact] + public async Task Should_forward_events_to_subscriber() + { + var events = Enumerable.Range(0, 50).Select(CreateEvent).ToArray(); + + var receivedEvents = new List(); + + A.CallTo(() => eventStore.QueryAllAsync("^my-stream", position, A._, A._)) + .Returns(events.ToAsyncEnumerable()); + + A.CallTo(() => eventSubscriber.OnNextAsync(A._, A._)) + .Invokes(x => receivedEvents.Add(x.GetArgument(1)!)); + + var sut = new PollingSubscription(eventStore, eventSubscriber, "^my-stream", position); + + sut.WakeUp(); + + await WaitAndStopAsync(sut); + + receivedEvents.Should().BeEquivalentTo(events); + } + + [Fact] + public async Task Should_receive_missing_events_with_second_pull() + { + var events1 = Enumerable.Range(0, 200).Where(x => x % 2 == 0).Select(CreateEvent).ToArray(); + var events2 = Enumerable.Range(0, 200).Where(x => x % 2 == 1).Select(CreateEvent).ToArray(); + + var receivedEvents = new List(); + + A.CallTo(() => eventStore.QueryAllAsync("^my-stream", position, A._, A._)) + .Returns(events1.ToAsyncEnumerable()); + + A.CallTo(() => eventStore.QueryAllAsync("^my-stream", "100", A._, A._)) + .Returns(events2.ToAsyncEnumerable()); + + A.CallTo(() => eventSubscriber.OnNextAsync(A._, A._)) + .Invokes(x => receivedEvents.Add(x.GetArgument(1)!)); + + var sut = new PollingSubscription(eventStore, eventSubscriber, "^my-stream", position); + + sut.WakeUp(); + + await WaitAndStopAsync(sut); + + receivedEvents.Should().BeEquivalentTo(events1.Union(events2)); + } + + private StoredEvent CreateEvent(int position) + { + return new StoredEvent( + "my-stream", + position.ToString(CultureInfo.InvariantCulture)!, + position, + new EventData( + "type", + new EnvelopeHeaders + { + [CommonHeaders.EventId] = Guid.NewGuid().ToString() + }, + "payload")); + } + private static async Task WaitAndStopAsync(IEventSubscription sut) { await Task.Delay(200);