Browse Source

Better polling subscription.

pull/906/head
Sebastian 4 years ago
parent
commit
c0d5d6f6b5
  1. 7
      backend/src/Squidex.Domain.Apps.Entities/Assets/AssetUsageTracker_EventHandling.cs
  2. 57
      backend/src/Squidex.Infrastructure/EventSourcing/PollingSubscription.cs
  3. 71
      backend/tests/Squidex.Infrastructure.Tests/EventSourcing/PollingSubscriptionTests.cs

7
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<DomainId>());
Console.WriteLine($"Writing tags: {JsonSerializer.Serialize(tagsPerApp, options)}");
await store.WriteManyAsync(tagsPerAsset.Select(x => new SnapshotWriteJob<State>(x.Key, x.Value, 0)));
}

57
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<Guid> eventIds = new HashSet<Guid>(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<StoredEvent> 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)
{

71
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<StoredEvent>();
A.CallTo(() => eventStore.QueryAllAsync("^my-stream", position, A<int>._, A<CancellationToken>._))
.Returns(events.ToAsyncEnumerable());
A.CallTo(() => eventSubscriber.OnNextAsync(A<IEventSubscription>._, A<StoredEvent>._))
.Invokes(x => receivedEvents.Add(x.GetArgument<StoredEvent>(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<StoredEvent>();
A.CallTo(() => eventStore.QueryAllAsync("^my-stream", position, A<int>._, A<CancellationToken>._))
.Returns(events1.ToAsyncEnumerable());
A.CallTo(() => eventStore.QueryAllAsync("^my-stream", "100", A<int>._, A<CancellationToken>._))
.Returns(events2.ToAsyncEnumerable());
A.CallTo(() => eventSubscriber.OnNextAsync(A<IEventSubscription>._, A<StoredEvent>._))
.Invokes(x => receivedEvents.Add(x.GetArgument<StoredEvent>(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);

Loading…
Cancel
Save