Browse Source

Event query improvements (#1033)

* Schema name.

* Fix build.

* Query improvements.

* Simplifications.

* Fix regex.

* Build fix.
pull/1034/head
Sebastian Stehle 3 years ago
committed by GitHub
parent
commit
5440e68bd8
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 25
      backend/src/Migrations/RebuilderExtensions.cs
  2. 2
      backend/src/Squidex.Domain.Apps.Core.Model/Schemas/ReferencesFieldProperties.cs
  3. 20
      backend/src/Squidex.Domain.Apps.Core.Operations/Subscriptions/SubscriptionPublisher.cs
  4. 20
      backend/src/Squidex.Domain.Apps.Entities.MongoDb/Schemas/MongoSchemasHash.cs
  5. 4
      backend/src/Squidex.Domain.Apps.Entities/Apps/AppEventDeleter.cs
  6. 10
      backend/src/Squidex.Domain.Apps.Entities/Apps/AppPermanentDeleter.cs
  7. 10
      backend/src/Squidex.Domain.Apps.Entities/Assets/AssetPermanentDeleter.cs
  8. 20
      backend/src/Squidex.Domain.Apps.Entities/Assets/AssetUsageTracker_EventHandling.cs
  9. 20
      backend/src/Squidex.Domain.Apps.Entities/Assets/AssetsFluidExtension.cs
  10. 4
      backend/src/Squidex.Domain.Apps.Entities/Assets/RebuildFiles.cs
  11. 10
      backend/src/Squidex.Domain.Apps.Entities/Assets/RecursiveDeleter.cs
  12. 9
      backend/src/Squidex.Domain.Apps.Entities/Backup/BackupProcessor.cs
  13. 2
      backend/src/Squidex.Domain.Apps.Entities/Comments/DomainObject/CommentsStream.cs
  14. 20
      backend/src/Squidex.Domain.Apps.Entities/Contents/ReferencesFluidExtension.cs
  15. 20
      backend/src/Squidex.Domain.Apps.Entities/Contents/Text/TextIndexingProcess.cs
  16. 10
      backend/src/Squidex.Domain.Apps.Entities/Invitation/InvitationEventConsumer.cs
  17. 5
      backend/src/Squidex.Domain.Apps.Entities/Rules/Runner/DefaultRuleRunnerService.cs
  18. 4
      backend/src/Squidex.Domain.Apps.Entities/Rules/Runner/RuleRunnerProcessor.cs
  19. 12
      backend/src/Squidex.Infrastructure.GetEventStore/EventSourcing/EventStoreProjectionClient.cs
  20. 32
      backend/src/Squidex.Infrastructure.GetEventStore/EventSourcing/GetEventStore.cs
  21. 4
      backend/src/Squidex.Infrastructure.GetEventStore/EventSourcing/GetEventStoreSubscription.cs
  22. 18
      backend/src/Squidex.Infrastructure.GetEventStore/EventSourcing/Utils.cs
  23. 45
      backend/src/Squidex.Infrastructure.MongoDb/EventSourcing/FilterExtensions.cs
  24. 10
      backend/src/Squidex.Infrastructure.MongoDb/EventSourcing/MongoEventStoreSubscription.cs
  25. 72
      backend/src/Squidex.Infrastructure.MongoDb/EventSourcing/MongoEventStore_Reader.cs
  26. 14
      backend/src/Squidex.Infrastructure.MongoDb/EventSourcing/MongoEventStore_Writer.cs
  27. 4
      backend/src/Squidex.Infrastructure/Commands/Rebuilder.cs
  28. 4
      backend/src/Squidex.Infrastructure/EventSourcing/IEventConsumer.cs
  29. 28
      backend/src/Squidex.Infrastructure/EventSourcing/IEventStore.cs
  30. 2
      backend/src/Squidex.Infrastructure/EventSourcing/PollingSubscription.cs
  31. 35
      backend/src/Squidex.Infrastructure/EventSourcing/StreamFilter.cs
  32. 14
      backend/src/Squidex.Infrastructure/EventSourcing/StreamFilterKind.cs
  33. 15
      backend/src/Squidex.Infrastructure/States/BatchContext.cs
  34. 4
      backend/src/Squidex.Infrastructure/States/Persistence.cs
  35. 2
      backend/tests/Squidex.Domain.Apps.Core.Tests/Operations/Subscriptions/SubscriptionPublisherTests.cs
  36. 4
      backend/tests/Squidex.Domain.Apps.Entities.Tests/Apps/AppEventDeleterTests.cs
  37. 4
      backend/tests/Squidex.Domain.Apps.Entities.Tests/Apps/AppPermanentDeleterTests.cs
  38. 4
      backend/tests/Squidex.Domain.Apps.Entities.Tests/Assets/AssetPermanentDeleterTests.cs
  39. 4
      backend/tests/Squidex.Domain.Apps.Entities.Tests/Assets/AssetUsageTrackerTests.cs
  40. 4
      backend/tests/Squidex.Domain.Apps.Entities.Tests/Assets/RecursiveDeleterTests.cs
  41. 4
      backend/tests/Squidex.Domain.Apps.Entities.Tests/Assets/RepairFilesTests.cs
  42. 6
      backend/tests/Squidex.Domain.Apps.Entities.Tests/Contents/Text/TextIndexerTestsBase.cs
  43. 19
      backend/tests/Squidex.Domain.Apps.Entities.Tests/Invitation/InvitationEventConsumerTests.cs
  44. 2
      backend/tests/Squidex.Domain.Apps.Entities.Tests/Rules/RuleEnqueuerTests.cs
  45. 22
      backend/tests/Squidex.Infrastructure.Tests/EventSourcing/Consume/EventConsumerProcessorTests.cs
  46. 62
      backend/tests/Squidex.Infrastructure.Tests/EventSourcing/EventStoreTests.cs
  47. 5
      backend/tests/Squidex.Infrastructure.Tests/EventSourcing/MongoEventStoreTests.cs
  48. 2
      backend/tests/Squidex.Infrastructure.Tests/EventSourcing/MongoParallelInsertTests.cs
  49. 2
      backend/tests/Squidex.Infrastructure.Tests/EventSourcing/PollingSubscriptionTests.cs
  50. 8
      backend/tests/Squidex.Infrastructure.Tests/EventSourcing/RetrySubscriptionTests.cs
  51. 27
      backend/tests/Squidex.Infrastructure.Tests/EventSourcing/StreamFilterTests.cs
  52. 15
      backend/tests/Squidex.Infrastructure.Tests/States/PersistenceBatchTests.cs
  53. 12
      backend/tests/Squidex.Infrastructure.Tests/States/PersistenceEventSourcingTests.cs
  54. 2
      backend/tests/Squidex.Infrastructure.Tests/States/PersistenceSnapshotTests.cs

25
backend/src/Migrations/RebuilderExtensions.cs

@ -11,6 +11,7 @@ using Squidex.Domain.Apps.Entities.Contents.DomainObject;
using Squidex.Domain.Apps.Entities.Rules.DomainObject; using Squidex.Domain.Apps.Entities.Rules.DomainObject;
using Squidex.Domain.Apps.Entities.Schemas.DomainObject; using Squidex.Domain.Apps.Entities.Schemas.DomainObject;
using Squidex.Infrastructure.Commands; using Squidex.Infrastructure.Commands;
using Squidex.Infrastructure.EventSourcing;
namespace Migrations; namespace Migrations;
@ -21,36 +22,48 @@ public static class RebuilderExtensions
public static Task RebuildAppsAsync(this Rebuilder rebuilder, int batchSize, public static Task RebuildAppsAsync(this Rebuilder rebuilder, int batchSize,
CancellationToken ct = default) CancellationToken ct = default)
{ {
return rebuilder.RebuildAsync<AppDomainObject, AppDomainObject.State>("^app\\-", batchSize, AllowedErrorRate, ct); var streamFilter = StreamFilter.Prefix("app-");
return rebuilder.RebuildAsync<AppDomainObject, AppDomainObject.State>(streamFilter, batchSize, AllowedErrorRate, ct);
} }
public static Task RebuildSchemasAsync(this Rebuilder rebuilder, int batchSize, public static Task RebuildSchemasAsync(this Rebuilder rebuilder, int batchSize,
CancellationToken ct = default) CancellationToken ct = default)
{ {
return rebuilder.RebuildAsync<SchemaDomainObject, SchemaDomainObject.State>("^schema\\-", batchSize, AllowedErrorRate, ct); var streamFilter = StreamFilter.Prefix("schema-");
return rebuilder.RebuildAsync<SchemaDomainObject, SchemaDomainObject.State>(streamFilter, batchSize, AllowedErrorRate, ct);
} }
public static Task RebuildRulesAsync(this Rebuilder rebuilder, int batchSize, public static Task RebuildRulesAsync(this Rebuilder rebuilder, int batchSize,
CancellationToken ct = default) CancellationToken ct = default)
{ {
return rebuilder.RebuildAsync<RuleDomainObject, RuleDomainObject.State>("^rule\\-", batchSize, AllowedErrorRate, ct); var streamFilter = StreamFilter.Prefix("rule-");
return rebuilder.RebuildAsync<RuleDomainObject, RuleDomainObject.State>(streamFilter, batchSize, AllowedErrorRate, ct);
} }
public static Task RebuildAssetsAsync(this Rebuilder rebuilder, int batchSize, public static Task RebuildAssetsAsync(this Rebuilder rebuilder, int batchSize,
CancellationToken ct = default) CancellationToken ct = default)
{ {
return rebuilder.RebuildAsync<AssetDomainObject, AssetDomainObject.State>("^asset\\-", batchSize, AllowedErrorRate, ct); var streamFilter = StreamFilter.Prefix("asset-");
return rebuilder.RebuildAsync<AssetDomainObject, AssetDomainObject.State>(streamFilter, batchSize, AllowedErrorRate, ct);
} }
public static Task RebuildAssetFoldersAsync(this Rebuilder rebuilder, int batchSize, public static Task RebuildAssetFoldersAsync(this Rebuilder rebuilder, int batchSize,
CancellationToken ct = default) CancellationToken ct = default)
{ {
return rebuilder.RebuildAsync<AssetFolderDomainObject, AssetFolderDomainObject.State>("^assetFolder\\-", batchSize, AllowedErrorRate, ct); var streamFilter = StreamFilter.Prefix("assetFolder-");
return rebuilder.RebuildAsync<AssetFolderDomainObject, AssetFolderDomainObject.State>(streamFilter, batchSize, AllowedErrorRate, ct);
} }
public static Task RebuildContentAsync(this Rebuilder rebuilder, int batchSize, public static Task RebuildContentAsync(this Rebuilder rebuilder, int batchSize,
CancellationToken ct = default) CancellationToken ct = default)
{ {
return rebuilder.RebuildAsync<ContentDomainObject, ContentDomainObject.State>("^content\\-", batchSize, AllowedErrorRate, ct); var streamFilter = StreamFilter.Prefix("content-");
return rebuilder.RebuildAsync<ContentDomainObject, ContentDomainObject.State>(streamFilter, batchSize, AllowedErrorRate, ct);
} }
} }

2
backend/src/Squidex.Domain.Apps.Core.Model/Schemas/ReferencesFieldProperties.cs

@ -26,7 +26,7 @@ public sealed record ReferencesFieldProperties : FieldProperties
public bool MustBePublished { get; init; } public bool MustBePublished { get; init; }
public string? Query { get; init;} public string? Query { get; init; }
public ReferencesFieldEditor Editor { get; init; } public ReferencesFieldEditor Editor { get; init; }

20
backend/src/Squidex.Domain.Apps.Core.Operations/Subscriptions/SubscriptionPublisher.cs

@ -16,25 +16,13 @@ public sealed class SubscriptionPublisher : IEventConsumer
private readonly ISubscriptionService subscriptionService; private readonly ISubscriptionService subscriptionService;
private readonly IEnumerable<ISubscriptionEventCreator> subscriptionEventCreators; private readonly IEnumerable<ISubscriptionEventCreator> subscriptionEventCreators;
public string Name public string Name => "Subscriptions";
{
get => "Subscriptions";
}
public string EventsFilter public StreamFilter EventsFilter { get; } = StreamFilter.Prefix("content-", "asset-");
{
get => "^(content-|asset-)";
}
public bool StartLatest public bool StartLatest => true;
{
get => true;
}
public bool CanClear public bool CanClear => false;
{
get => false;
}
public SubscriptionPublisher(ISubscriptionService subscriptionService, IEnumerable<ISubscriptionEventCreator> subscriptionEventCreators) public SubscriptionPublisher(ISubscriptionService subscriptionService, IEnumerable<ISubscriptionEventCreator> subscriptionEventCreators)
{ {

20
backend/src/Squidex.Domain.Apps.Entities.MongoDb/Schemas/MongoSchemasHash.cs

@ -19,25 +19,11 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Schemas;
public sealed class MongoSchemasHash : MongoRepositoryBase<MongoSchemasHashEntity>, ISchemasHash, IEventConsumer, IDeleter public sealed class MongoSchemasHash : MongoRepositoryBase<MongoSchemasHashEntity>, ISchemasHash, IEventConsumer, IDeleter
{ {
public int BatchSize public int BatchSize => 1000;
{
get => 1000;
}
public int BatchDelay public int BatchDelay => 500;
{
get => 500;
}
public string Name public StreamFilter EventsFilter { get; } = StreamFilter.Prefix("schema-");
{
get => GetType().Name;
}
public string EventsFilter
{
get => "^schema-";
}
public MongoSchemasHash(IMongoDatabase database) public MongoSchemasHash(IMongoDatabase database)
: base(database) : base(database)

4
backend/src/Squidex.Domain.Apps.Entities/Apps/AppEventDeleter.cs

@ -23,6 +23,8 @@ public sealed class AppEventDeleter : IDeleter
public Task DeleteAppAsync(IAppEntity app, public Task DeleteAppAsync(IAppEntity app,
CancellationToken ct) CancellationToken ct)
{ {
return eventStore.DeleteAsync($"^([a-zA-Z0-9]+)\\-{app.Id}", ct); var streamFilter = StreamFilter.Prefix($"([a-zA-Z0-9]+)-{app.Id}");
return eventStore.DeleteAsync(streamFilter, ct);
} }
} }

10
backend/src/Squidex.Domain.Apps.Entities/Apps/AppPermanentDeleter.cs

@ -20,15 +20,7 @@ public sealed class AppPermanentDeleter : IEventConsumer
private readonly IDomainObjectFactory factory; private readonly IDomainObjectFactory factory;
private readonly HashSet<string> consumingTypes; private readonly HashSet<string> consumingTypes;
public string Name public StreamFilter EventsFilter { get; } = StreamFilter.Prefix("app-");
{
get => GetType().Name;
}
public string EventsFilter
{
get => "^app-";
}
public AppPermanentDeleter(IEnumerable<IDeleter> deleters, IDomainObjectFactory factory, TypeRegistry typeRegistry) public AppPermanentDeleter(IEnumerable<IDeleter> deleters, IDomainObjectFactory factory, TypeRegistry typeRegistry)
{ {

10
backend/src/Squidex.Domain.Apps.Entities/Assets/AssetPermanentDeleter.cs

@ -17,15 +17,7 @@ public sealed class AssetPermanentDeleter : IEventConsumer
private readonly IAssetFileStore assetFileStore; private readonly IAssetFileStore assetFileStore;
private readonly HashSet<string> consumingTypes; private readonly HashSet<string> consumingTypes;
public string Name public StreamFilter EventsFilter { get; } = StreamFilter.Prefix("asset-");
{
get => GetType().Name;
}
public string EventsFilter
{
get => "^asset-";
}
public AssetPermanentDeleter(IAssetFileStore assetFileStore, TypeRegistry typeRegistry) public AssetPermanentDeleter(IAssetFileStore assetFileStore, TypeRegistry typeRegistry)
{ {

20
backend/src/Squidex.Domain.Apps.Entities/Assets/AssetUsageTracker_EventHandling.cs

@ -21,25 +21,11 @@ public partial class AssetUsageTracker : IEventConsumer
{ {
private IMemoryCache memoryCache; private IMemoryCache memoryCache;
public int BatchSize public int BatchSize => 1000;
{
get => 1000;
}
public int BatchDelay public int BatchDelay => 1000;
{
get => 1000;
}
public string Name public StreamFilter EventsFilter { get; } = StreamFilter.Prefix("asset-");
{
get => GetType().Name;
}
public string EventsFilter
{
get => "^asset-";
}
private void ClearCache() private void ClearCache()
{ {

20
backend/src/Squidex.Domain.Apps.Entities/Assets/AssetsFluidExtension.cs

@ -39,19 +39,21 @@ public sealed class AssetsFluidExtension : IFluidExtension
private async ValueTask<Completion> ResolveAsset(ValueTuple<Expression, Expression> arguments, TextWriter writer, TextEncoder encoder, TemplateContext context) private async ValueTask<Completion> ResolveAsset(ValueTuple<Expression, Expression> arguments, TextWriter writer, TextEncoder encoder, TemplateContext context)
{ {
if (context.GetValue("event")?.ToObjectValue() is EnrichedEvent enrichedEvent) if (context.GetValue("event")?.ToObjectValue() is not EnrichedEvent enrichedEvent)
{ {
var (nameArg, idArg) = arguments; return Completion.Normal;
}
var assetId = await idArg.EvaluateAsync(context); var (nameArg, idArg) = arguments;
var asset = await ResolveAssetAsync(serviceProvider, enrichedEvent.AppId.Id, assetId);
if (asset != null) var assetId = await idArg.EvaluateAsync(context);
{ var asset = await ResolveAssetAsync(serviceProvider, enrichedEvent.AppId.Id, assetId);
var name = (await nameArg.EvaluateAsync(context)).ToStringValue();
context.SetValue(name, asset); if (asset != null)
} {
var name = (await nameArg.EvaluateAsync(context)).ToStringValue();
context.SetValue(name, asset);
} }
return Completion.Normal; return Completion.Normal;

4
backend/src/Squidex.Domain.Apps.Entities/Assets/RebuildFiles.cs

@ -33,7 +33,9 @@ public sealed class RebuildFiles
public async Task RepairAsync( public async Task RepairAsync(
CancellationToken ct = default) CancellationToken ct = default)
{ {
await foreach (var storedEvent in eventStore.QueryAllAsync("^asset\\-", ct: ct)) var streamFilter = StreamFilter.Prefix("asset-");
await foreach (var storedEvent in eventStore.QueryAllAsync(streamFilter, ct: ct))
{ {
var @event = eventFormatter.ParseIfKnown(storedEvent); var @event = eventFormatter.ParseIfKnown(storedEvent);

10
backend/src/Squidex.Domain.Apps.Entities/Assets/RecursiveDeleter.cs

@ -23,15 +23,7 @@ public sealed class RecursiveDeleter : IEventConsumer
private readonly ILogger<RecursiveDeleter> log; private readonly ILogger<RecursiveDeleter> log;
private readonly HashSet<string> consumingTypes; private readonly HashSet<string> consumingTypes;
public string Name public StreamFilter EventsFilter { get; } = StreamFilter.Prefix("assetFolder-");
{
get => GetType().Name;
}
public string EventsFilter
{
get => "^assetFolder-";
}
public RecursiveDeleter( public RecursiveDeleter(
ICommandBus commandBus, ICommandBus commandBus,

9
backend/src/Squidex.Domain.Apps.Entities/Backup/BackupProcessor.cs

@ -147,7 +147,9 @@ public sealed partial class BackupProcessor
var backupUsers = new UserMapping(run.Actor); var backupUsers = new UserMapping(run.Actor);
var backupContext = new BackupContext(appId, backupUsers, writer); var backupContext = new BackupContext(appId, backupUsers, writer);
await foreach (var storedEvent in eventStore.QueryAllAsync(GetFilter(), ct: ct)) var streamFilter = StreamFilter.Prefix($"[^\\-]*-{appId}");
await foreach (var storedEvent in eventStore.QueryAllAsync(streamFilter, ct: ct))
{ {
var @event = eventFormatter.Parse(storedEvent); var @event = eventFormatter.Parse(storedEvent);
@ -200,11 +202,6 @@ public sealed partial class BackupProcessor
} }
} }
private string GetFilter()
{
return $"^[^\\-]*-{Regex.Escape(appId.ToString())}";
}
public Task DeleteAsync(DomainId id) public Task DeleteAsync(DomainId id)
{ {
return scheduler.ScheduleAsync(async _ => return scheduler.ScheduleAsync(async _ =>

2
backend/src/Squidex.Domain.Apps.Entities/Comments/DomainObject/CommentsStream.cs

@ -42,7 +42,7 @@ public class CommentsStream : IAggregate
public virtual async Task LoadAsync( public virtual async Task LoadAsync(
CancellationToken ct) CancellationToken ct)
{ {
var storedEvents = await eventStore.QueryReverseAsync(streamName, 100, ct); var storedEvents = await eventStore.QueryStreamReverseAsync(streamName, 100, ct);
foreach (var @event in storedEvents) foreach (var @event in storedEvents)
{ {

20
backend/src/Squidex.Domain.Apps.Entities/Contents/ReferencesFluidExtension.cs

@ -39,19 +39,21 @@ public sealed class ReferencesFluidExtension : IFluidExtension
private async ValueTask<Completion> ResolveReference(ValueTuple<Expression, Expression> arguments, TextWriter writer, TextEncoder encoder, TemplateContext context) private async ValueTask<Completion> ResolveReference(ValueTuple<Expression, Expression> arguments, TextWriter writer, TextEncoder encoder, TemplateContext context)
{ {
if (context.GetValue("event")?.ToObjectValue() is EnrichedEvent enrichedEvent) if (context.GetValue("event")?.ToObjectValue() is not EnrichedEvent enrichedEvent)
{ {
var (nameArg, idArg) = arguments; return Completion.Normal;
}
var contentId = await idArg.EvaluateAsync(context); var (nameArg, idArg) = arguments;
var content = await ResolveContentAsync(serviceProvider, enrichedEvent.AppId.Id, contentId);
if (content != null) var contentId = await idArg.EvaluateAsync(context);
{ var content = await ResolveContentAsync(serviceProvider, enrichedEvent.AppId.Id, contentId);
var name = (await nameArg.EvaluateAsync(context)).ToStringValue();
context.SetValue(name, content); if (content != null)
} {
var name = (await nameArg.EvaluateAsync(context)).ToStringValue();
context.SetValue(name, content);
} }
return Completion.Normal; return Completion.Normal;

20
backend/src/Squidex.Domain.Apps.Entities/Contents/Text/TextIndexingProcess.cs

@ -21,25 +21,13 @@ public sealed class TextIndexingProcess : IEventConsumer
private readonly ITextIndex textIndex; private readonly ITextIndex textIndex;
private readonly ITextIndexerState textIndexerState; private readonly ITextIndexerState textIndexerState;
public int BatchSize public int BatchSize => 1000;
{
get => 1000;
}
public int BatchDelay public int BatchDelay => 1000;
{
get => 1000;
}
public string Name public string Name => "TextIndexer5";
{
get => "TextIndexer5";
}
public string EventsFilter public StreamFilter EventsFilter { get; } = StreamFilter.Prefix("content-");
{
get => "^content-";
}
public ITextIndex TextIndex public ITextIndex TextIndex
{ {

10
backend/src/Squidex.Domain.Apps.Entities/Invitation/InvitationEventConsumer.cs

@ -24,15 +24,9 @@ public sealed class InvitationEventConsumer : IEventConsumer
private readonly IAppProvider appProvider; private readonly IAppProvider appProvider;
private readonly ILogger<InvitationEventConsumer> log; private readonly ILogger<InvitationEventConsumer> log;
public string Name public string Name => "NotificationEmailSender";
{
get => "NotificationEmailSender";
}
public string EventsFilter public StreamFilter EventsFilter { get; } = StreamFilter.Prefix("app-");
{
get { return "^app-|^app-"; }
}
public InvitationEventConsumer( public InvitationEventConsumer(
IAppProvider appProvider, IAppProvider appProvider,

5
backend/src/Squidex.Domain.Apps.Entities/Rules/Runner/DefaultRuleRunnerService.cs

@ -65,9 +65,10 @@ public sealed class DefaultRuleRunnerService : IRuleRunnerService
var simulatedEvents = new List<SimulatedRuleEvent>(MaxSimulatedEvents); var simulatedEvents = new List<SimulatedRuleEvent>(MaxSimulatedEvents);
var fromNow = SystemClock.Instance.GetCurrentInstant().Minus(Duration.FromDays(7)); var streamStart = SystemClock.Instance.GetCurrentInstant().Minus(Duration.FromDays(7));
var streamFilter = StreamFilter.Prefix($"([a-zA-Z0-9]+)-{appId.Id}");
await foreach (var storedEvent in eventStore.QueryAllReverseAsync($"^([a-zA-Z0-9]+)\\-{appId.Id}", fromNow, MaxSimulatedEvents, ct)) await foreach (var storedEvent in eventStore.QueryAllReverseAsync(streamFilter, streamStart, MaxSimulatedEvents, ct))
{ {
var @event = eventFormatter.ParseIfKnown(storedEvent); var @event = eventFormatter.ParseIfKnown(storedEvent);

4
backend/src/Squidex.Domain.Apps.Entities/Rules/Runner/RuleRunnerProcessor.cs

@ -254,9 +254,9 @@ public sealed class RuleRunnerProcessor
await using var batch = new RuleQueueWriter(ruleEventRepository, ruleUsageTracker, null); await using var batch = new RuleQueueWriter(ruleEventRepository, ruleUsageTracker, null);
// Use a prefix query so that the storage can use an index for the query. // Use a prefix query so that the storage can use an index for the query.
var filter = $"^([a-z]+)\\-{appId}"; var streamFilter = StreamFilter.Prefix($"([a-zA-Z0-9]+)\\-{appId}");
await foreach (var storedEvent in eventStore.QueryAllAsync(filter, run.Job.Position, ct: ct)) await foreach (var storedEvent in eventStore.QueryAllAsync(streamFilter, run.Job.Position, ct: ct))
{ {
var @event = eventFormatter.ParseIfKnown(storedEvent); var @event = eventFormatter.ParseIfKnown(storedEvent);

12
backend/src/Squidex.Infrastructure.GetEventStore/EventSourcing/EventStoreProjectionClient.cs

@ -29,22 +29,22 @@ public sealed class EventStoreProjectionClient
return $"by-{projectionPrefix.Slugify()}-{filter.Slugify()}"; return $"by-{projectionPrefix.Slugify()}-{filter.Slugify()}";
} }
public async Task<string> CreateProjectionAsync(string? streamFilter = null) public async Task<string> CreateProjectionAsync(StreamFilter filter)
{ {
if (!string.IsNullOrWhiteSpace(streamFilter) && streamFilter[0] != '^') if (filter.Kind == StreamFilterKind.MatchFull && filter.Prefixes?.Count == 1)
{ {
return $"{projectionPrefix}-{streamFilter}"; return $"{projectionPrefix}-{filter.Prefixes[0]}";
} }
streamFilter ??= ".*"; var regex = filter.ToRegex();
var name = CreateFilterProjectionName(streamFilter); var name = CreateFilterProjectionName(regex);
var query = var query =
$@"fromAll() $@"fromAll()
.when({{ .when({{
$any: function (s, e) {{ $any: function (s, e) {{
if (e.streamId.indexOf('{projectionPrefix}') === 0 && /{streamFilter}/.test(e.streamId.substring({projectionPrefix.Length + 1}))) {{ if (e.streamId.indexOf('{projectionPrefix}') === 0 && /{regex}/.test(e.streamId.substring({projectionPrefix.Length + 1}))) {{
linkTo('{name}', e); linkTo('{name}', e);
}} }}
}} }}

32
backend/src/Squidex.Infrastructure.GetEventStore/EventSourcing/GetEventStore.cs

@ -46,14 +46,14 @@ public sealed class GetEventStore : IEventStore, IInitializable
} }
} }
public IEventSubscription CreateSubscription(IEventSubscriber<StoredEvent> subscriber, string? streamFilter = null, string? position = null) public IEventSubscription CreateSubscription(IEventSubscriber<StoredEvent> subscriber, StreamFilter filter, string? position = null)
{ {
Guard.NotNull(streamFilter); Guard.NotNull(filter);
return new GetEventStoreSubscription(subscriber, client, projectionClient, serializer, position, StreamPrefix, streamFilter); return new GetEventStoreSubscription(subscriber, client, projectionClient, serializer, position, StreamPrefix, filter);
} }
public async IAsyncEnumerable<StoredEvent> QueryAllAsync(string? streamFilter = null, string? position = null, int take = int.MaxValue, public async IAsyncEnumerable<StoredEvent> QueryAllAsync(StreamFilter filter, string? position = null, int take = int.MaxValue,
[EnumeratorCancellation] CancellationToken ct = default) [EnumeratorCancellation] CancellationToken ct = default)
{ {
if (take <= 0) if (take <= 0)
@ -61,7 +61,7 @@ public sealed class GetEventStore : IEventStore, IInitializable
yield break; yield break;
} }
var streamName = await projectionClient.CreateProjectionAsync(streamFilter); var streamName = await projectionClient.CreateProjectionAsync(filter);
var stream = QueryAsync(streamName, position.ToPosition(false), take, ct); var stream = QueryAsync(streamName, position.ToPosition(false), take, ct);
@ -71,7 +71,7 @@ public sealed class GetEventStore : IEventStore, IInitializable
} }
} }
public async IAsyncEnumerable<StoredEvent> QueryAllReverseAsync(string? streamFilter = null, Instant timestamp = default, int take = int.MaxValue, public async IAsyncEnumerable<StoredEvent> QueryAllReverseAsync(StreamFilter filter, Instant timestamp = default, int take = int.MaxValue,
[EnumeratorCancellation] CancellationToken ct = default) [EnumeratorCancellation] CancellationToken ct = default)
{ {
if (take <= 0) if (take <= 0)
@ -79,7 +79,7 @@ public sealed class GetEventStore : IEventStore, IInitializable
yield break; yield break;
} }
var streamName = await projectionClient.CreateProjectionAsync(streamFilter); var streamName = await projectionClient.CreateProjectionAsync(filter);
var stream = QueryReverseAsync(streamName, StreamPosition.End, take, ct); var stream = QueryReverseAsync(streamName, StreamPosition.End, take, ct);
@ -89,11 +89,9 @@ public sealed class GetEventStore : IEventStore, IInitializable
} }
} }
public async Task<IReadOnlyList<StoredEvent>> QueryReverseAsync(string streamName, int count = int.MaxValue, public async Task<IReadOnlyList<StoredEvent>> QueryStreamReverseAsync(string streamName, int count = int.MaxValue,
CancellationToken ct = default) CancellationToken ct = default)
{ {
Guard.NotNullOrEmpty(streamName);
if (count <= 0) if (count <= 0)
{ {
return EmptyEvents; return EmptyEvents;
@ -103,7 +101,7 @@ public sealed class GetEventStore : IEventStore, IInitializable
{ {
var result = new List<StoredEvent>(); var result = new List<StoredEvent>();
var stream = QueryReverseAsync(GetStreamName(streamName), StreamPosition.End, count, ct); var stream = QueryReverseAsync(streamName, StreamPosition.End, count, ct);
await foreach (var storedEvent in stream.IgnoreNotFound(ct)) await foreach (var storedEvent in stream.IgnoreNotFound(ct))
{ {
@ -114,16 +112,14 @@ public sealed class GetEventStore : IEventStore, IInitializable
} }
} }
public async Task<IReadOnlyList<StoredEvent>> QueryAsync(string streamName, long afterStreamPosition = EtagVersion.Empty, public async Task<IReadOnlyList<StoredEvent>> QueryStreamAsync(string streamName, long afterStreamPosition = EtagVersion.Empty,
CancellationToken ct = default) CancellationToken ct = default)
{ {
Guard.NotNullOrEmpty(streamName);
using (Telemetry.Activities.StartActivity("GetEventStore/QueryAsync")) using (Telemetry.Activities.StartActivity("GetEventStore/QueryAsync"))
{ {
var result = new List<StoredEvent>(); var result = new List<StoredEvent>();
var stream = QueryAsync(GetStreamName(streamName), afterStreamPosition.ToPositionBefore(), int.MaxValue, ct); var stream = QueryAsync(streamName, afterStreamPosition.ToPositionBefore(), int.MaxValue, ct);
await foreach (var storedEvent in stream.IgnoreNotFound(ct)) await foreach (var storedEvent in stream.IgnoreNotFound(ct))
{ {
@ -156,7 +152,7 @@ public sealed class GetEventStore : IEventStore, IInitializable
streamName, streamName,
start, start,
count, count,
resolveLinkTos: true, true,
cancellationToken: ct); cancellationToken: ct);
return result.Select(x => Formatter.Read(x, StreamPrefix, serializer)); return result.Select(x => Formatter.Read(x, StreamPrefix, serializer));
@ -210,10 +206,10 @@ public sealed class GetEventStore : IEventStore, IInitializable
} }
} }
public async Task DeleteAsync(string streamFilter, public async Task DeleteAsync(StreamFilter filter,
CancellationToken ct = default) CancellationToken ct = default)
{ {
var streamName = await projectionClient.CreateProjectionAsync(streamFilter); var streamName = await projectionClient.CreateProjectionAsync(filter);
var events = client.ReadStreamAsync(Direction.Forwards, streamName, StreamPosition.Start, resolveLinkTos: true, cancellationToken: ct); var events = client.ReadStreamAsync(Direction.Forwards, streamName, StreamPosition.Start, resolveLinkTos: true, cancellationToken: ct);

4
backend/src/Squidex.Infrastructure.GetEventStore/EventSourcing/GetEventStoreSubscription.cs

@ -23,14 +23,14 @@ internal sealed class GetEventStoreSubscription : IEventSubscription
IJsonSerializer serializer, IJsonSerializer serializer,
string? position, string? position,
string? prefix, string? prefix,
string? streamFilter) StreamFilter filter)
{ {
#pragma warning disable MA0134 // Observe result of async calls #pragma warning disable MA0134 // Observe result of async calls
Task.Run(async () => Task.Run(async () =>
{ {
var ct = cts.Token; var ct = cts.Token;
var streamName = await projectionClient.CreateProjectionAsync(streamFilter); var streamName = await projectionClient.CreateProjectionAsync(filter);
async Task OnEvent(StreamSubscription subscription, ResolvedEvent @event, async Task OnEvent(StreamSubscription subscription, ResolvedEvent @event,
CancellationToken ct) CancellationToken ct)

18
backend/src/Squidex.Infrastructure.GetEventStore/EventSourcing/Utils.cs

@ -7,6 +7,7 @@
using System.Globalization; using System.Globalization;
using System.Runtime.CompilerServices; using System.Runtime.CompilerServices;
using System.Text.RegularExpressions;
using EventStore.Client; using EventStore.Client;
namespace Squidex.Infrastructure.EventSourcing; namespace Squidex.Infrastructure.EventSourcing;
@ -77,4 +78,21 @@ public static class Utils
yield return enumerator.Current; yield return enumerator.Current;
} }
} }
public static string ToRegex(this StreamFilter filter)
{
if (filter.Prefixes == null)
{
return ".*";
}
if (filter.Kind == StreamFilterKind.MatchStart)
{
return $"^{string.Join('|', filter.Prefixes.Select(p => $"({p})"))}";
}
else
{
return $"^{string.Join('|', filter.Prefixes.Select(p => $"({p})"))}$";
}
}
} }

45
backend/src/Squidex.Infrastructure.MongoDb/EventSourcing/FilterExtensions.cs

@ -5,6 +5,7 @@
// All rights reserved. Licensed under the MIT license. // All rights reserved. Licensed under the MIT license.
// ========================================================================== // ==========================================================================
using System.Text.RegularExpressions;
using MongoDB.Driver; using MongoDB.Driver;
namespace Squidex.Infrastructure.EventSourcing; namespace Squidex.Infrastructure.EventSourcing;
@ -13,53 +14,57 @@ internal static class FilterExtensions
{ {
public static FilterDefinition<MongoEventCommit> ByOffset(long streamPosition) public static FilterDefinition<MongoEventCommit> ByOffset(long streamPosition)
{ {
return Builders<MongoEventCommit>.Filter.Gte(x => x.EventStreamOffset, streamPosition); var builder = Builders<MongoEventCommit>.Filter;
return builder.Gte(x => x.EventStreamOffset, streamPosition);
} }
public static FilterDefinition<MongoEventCommit> ByPosition(StreamPosition streamPosition) public static FilterDefinition<MongoEventCommit> ByPosition(StreamPosition streamPosition)
{ {
var builder = Builders<MongoEventCommit>.Filter;
if (streamPosition.IsEndOfCommit) if (streamPosition.IsEndOfCommit)
{ {
return Builders<MongoEventCommit>.Filter.Gt(x => x.Timestamp, streamPosition.Timestamp); return builder.Gt(x => x.Timestamp, streamPosition.Timestamp);
} }
else else
{ {
return Builders<MongoEventCommit>.Filter.Gte(x => x.Timestamp, streamPosition.Timestamp); return builder.Gte(x => x.Timestamp, streamPosition.Timestamp);
} }
} }
public static FilterDefinition<MongoEventCommit>? ByStream(string? streamFilter) public static FilterDefinition<MongoEventCommit> ByStream(StreamFilter filter)
{ {
if (StreamFilter.IsAll(streamFilter)) var builder = Builders<MongoEventCommit>.Filter;
{
return Builders<MongoEventCommit>.Filter.Exists(x => x.EventStream, true);
}
if (streamFilter.Contains('^', StringComparison.Ordinal)) if (filter.Prefixes == null)
{ {
return Builders<MongoEventCommit>.Filter.Regex(x => x.EventStream, streamFilter); return builder.Exists(x => x.EventStream, true);
} }
else
if (filter.Kind == StreamFilterKind.MatchStart)
{ {
return Builders<MongoEventCommit>.Filter.Eq(x => x.EventStream, streamFilter); return builder.Or(filter.Prefixes.Select(p => builder.Regex(x => x.EventStream, $"^{p}")));
} }
return builder.In(x => x.EventStream, filter.Prefixes);
} }
public static FilterDefinition<ChangeStreamDocument<MongoEventCommit>>? ByChangeInStream(string? streamFilter) public static FilterDefinition<ChangeStreamDocument<MongoEventCommit>>? ByChangeInStream(StreamFilter filter)
{ {
if (StreamFilter.IsAll(streamFilter)) var builder = Builders<ChangeStreamDocument<MongoEventCommit>>.Filter;
if (filter.Prefixes == null)
{ {
return null; return null;
} }
if (streamFilter.Contains('^', StringComparison.Ordinal)) if (filter.Kind == StreamFilterKind.MatchStart)
{
return Builders<ChangeStreamDocument<MongoEventCommit>>.Filter.Regex(x => x.FullDocument.EventStream, streamFilter);
}
else
{ {
return Builders<ChangeStreamDocument<MongoEventCommit>>.Filter.Eq(x => x.FullDocument.EventStream, streamFilter); return builder.Or(filter.Prefixes.Select(p => builder.Regex(x => x.FullDocument.EventStream, $"^{Regex.Escape(p)}")));
} }
return builder.In(x => x.FullDocument.EventStream, filter.Prefixes);
} }
public static IEnumerable<StoredEvent> Filtered(this MongoEventCommit commit, StreamPosition position) public static IEnumerable<StoredEvent> Filtered(this MongoEventCommit commit, StreamPosition position)

10
backend/src/Squidex.Infrastructure.MongoDb/EventSourcing/MongoEventStoreSubscription.cs

@ -18,7 +18,7 @@ public sealed class MongoEventStoreSubscription : IEventSubscription
private readonly IEventSubscriber<StoredEvent> eventSubscriber; private readonly IEventSubscriber<StoredEvent> eventSubscriber;
private readonly CancellationTokenSource stopToken = new CancellationTokenSource(); private readonly CancellationTokenSource stopToken = new CancellationTokenSource();
public MongoEventStoreSubscription(MongoEventStore eventStore, IEventSubscriber<StoredEvent> eventSubscriber, string? streamFilter, string? position) public MongoEventStoreSubscription(MongoEventStore eventStore, IEventSubscriber<StoredEvent> eventSubscriber, StreamFilter streamFilter, string? position)
{ {
this.eventStore = eventStore; this.eventStore = eventStore;
this.eventSubscriber = eventSubscriber; this.eventSubscriber = eventSubscriber;
@ -26,7 +26,7 @@ public sealed class MongoEventStoreSubscription : IEventSubscription
QueryAsync(streamFilter, position).Forget(); QueryAsync(streamFilter, position).Forget();
} }
private async Task QueryAsync(string? streamFilter, string? position) private async Task QueryAsync(StreamFilter streamFilter, string? position)
{ {
try try
{ {
@ -51,7 +51,7 @@ public sealed class MongoEventStoreSubscription : IEventSubscription
} }
} }
private async Task QueryCurrentAsync(string? streamFilter, StreamPosition lastPosition) private async Task QueryCurrentAsync(StreamFilter streamFilter, StreamPosition lastPosition)
{ {
BsonDocument? resumeToken = null; BsonDocument? resumeToken = null;
@ -103,7 +103,7 @@ public sealed class MongoEventStoreSubscription : IEventSubscription
} }
} }
private async Task<string?> QueryOldAsync(string? streamFilter, string? position) private async Task<string?> QueryOldAsync(StreamFilter streamFilter, string? position)
{ {
string? lastRawPosition = null; string? lastRawPosition = null;
@ -134,7 +134,7 @@ public sealed class MongoEventStoreSubscription : IEventSubscription
return lastRawPosition; return lastRawPosition;
} }
private static PipelineDefinition<ChangeStreamDocument<MongoEventCommit>, ChangeStreamDocument<MongoEventCommit>>? Match(string? streamFilter) private static PipelineDefinition<ChangeStreamDocument<MongoEventCommit>, ChangeStreamDocument<MongoEventCommit>>? Match(StreamFilter streamFilter)
{ {
var result = new EmptyPipelineDefinition<ChangeStreamDocument<MongoEventCommit>>(); var result = new EmptyPipelineDefinition<ChangeStreamDocument<MongoEventCommit>>();

72
backend/src/Squidex.Infrastructure.MongoDb/EventSourcing/MongoEventStore_Reader.cs

@ -20,25 +20,23 @@ public partial class MongoEventStore : MongoRepositoryBase<MongoEventCommit>, IE
{ {
private static readonly List<StoredEvent> EmptyEvents = new List<StoredEvent>(); private static readonly List<StoredEvent> EmptyEvents = new List<StoredEvent>();
public IEventSubscription CreateSubscription(IEventSubscriber<StoredEvent> subscriber, string? streamFilter = null, string? position = null) public IEventSubscription CreateSubscription(IEventSubscriber<StoredEvent> subscriber, StreamFilter filter, string? position = null)
{ {
Guard.NotNull(subscriber); Guard.NotNull(subscriber);
if (CanUseChangeStreams) if (CanUseChangeStreams)
{ {
return new MongoEventStoreSubscription(this, subscriber, streamFilter, position); return new MongoEventStoreSubscription(this, subscriber, filter, position);
} }
else else
{ {
return new PollingSubscription(this, subscriber, streamFilter, position); return new PollingSubscription(this, subscriber, filter, position);
} }
} }
public async Task<IReadOnlyList<StoredEvent>> QueryReverseAsync(string streamName, int count = int.MaxValue, public async Task<IReadOnlyList<StoredEvent>> QueryStreamReverseAsync(string streamName, int count = int.MaxValue,
CancellationToken ct = default) CancellationToken ct = default)
{ {
Guard.NotNullOrEmpty(streamName);
if (count <= 0) if (count <= 0)
{ {
return EmptyEvents; return EmptyEvents;
@ -46,7 +44,7 @@ public partial class MongoEventStore : MongoRepositoryBase<MongoEventCommit>, IE
using (Telemetry.Activities.StartActivity("MongoEventStore/QueryLatestAsync")) using (Telemetry.Activities.StartActivity("MongoEventStore/QueryLatestAsync"))
{ {
var filter = Filter.Eq(x => x.EventStream, streamName); var filter = FilterExtensions.ByStream(StreamFilter.Name(streamName));
var commits = var commits =
await Collection.Find(filter).Sort(Sort.Descending(x => x.Timestamp)).Limit(count) await Collection.Find(filter).Sort(Sort.Descending(x => x.Timestamp)).Limit(count)
@ -58,33 +56,26 @@ public partial class MongoEventStore : MongoRepositoryBase<MongoEventCommit>, IE
} }
} }
public async Task<IReadOnlyList<StoredEvent>> QueryAsync(string streamName, long afterStreamPosition = EtagVersion.Empty, public async Task<IReadOnlyList<StoredEvent>> QueryStreamAsync(string streamName, long afterStreamPosition = EtagVersion.Empty,
CancellationToken ct = default) CancellationToken ct = default)
{ {
Guard.NotNullOrEmpty(streamName);
using (Telemetry.Activities.StartActivity("MongoEventStore/QueryAsync")) using (Telemetry.Activities.StartActivity("MongoEventStore/QueryAsync"))
{ {
var filter =
Filter.And(
Filter.Eq(x => x.EventStream, streamName),
Filter.Gte(x => x.EventStreamOffset, afterStreamPosition));
var commits = var commits =
await Collection.Find(filter) await Collection.Find(CreateFilter(StreamFilter.Name(streamName), afterStreamPosition))
.ToListAsync(ct); .ToListAsync(ct);
var result = Convert(commits, afterStreamPosition); var result = Convert(commits, afterStreamPosition);
if ((commits.Count == 0 || commits[0].EventStreamOffset != afterStreamPosition) && afterStreamPosition > EtagVersion.Empty) if ((commits.Count == 0 || commits[0].EventStreamOffset != afterStreamPosition) && afterStreamPosition > EtagVersion.Empty)
{ {
filter = var filterBefore =
Filter.And( Filter.And(
Filter.Eq(x => x.EventStream, streamName), FilterExtensions.ByStream(StreamFilter.Name(streamName)),
Filter.Lt(x => x.EventStreamOffset, afterStreamPosition)); Filter.Lt(x => x.EventStreamOffset, afterStreamPosition));
commits = commits =
await Collection.Find(filter).SortByDescending(x => x.EventStreamOffset).Limit(1) await Collection.Find(filterBefore).SortByDescending(x => x.EventStreamOffset).Limit(1)
.ToListAsync(ct); .ToListAsync(ct);
result = Convert(commits, afterStreamPosition).Concat(result).ToList(); result = Convert(commits, afterStreamPosition).Concat(result).ToList();
@ -94,26 +85,7 @@ public partial class MongoEventStore : MongoRepositoryBase<MongoEventCommit>, IE
} }
} }
public async Task<IReadOnlyDictionary<string, IReadOnlyList<StoredEvent>>> QueryManyAsync(IEnumerable<string> streamNames, public async IAsyncEnumerable<StoredEvent> QueryAllReverseAsync(StreamFilter filter, Instant timestamp = default, int take = int.MaxValue,
CancellationToken ct = default)
{
Guard.NotNull(streamNames);
using (Telemetry.Activities.StartActivity("MongoEventStore/QueryManyAsync"))
{
var filter = Filter.In(x => x.EventStream, streamNames);
var commits =
await Collection.Find(filter)
.ToListAsync(ct);
var result = commits.GroupBy(x => x.EventStream).ToDictionary(x => x.Key, Convert);
return result;
}
}
public async IAsyncEnumerable<StoredEvent> QueryAllReverseAsync(string? streamFilter = null, Instant timestamp = default, int take = int.MaxValue,
[EnumeratorCancellation] CancellationToken ct = default) [EnumeratorCancellation] CancellationToken ct = default)
{ {
if (take <= 0) if (take <= 0)
@ -123,10 +95,8 @@ public partial class MongoEventStore : MongoRepositoryBase<MongoEventCommit>, IE
StreamPosition lastPosition = timestamp; StreamPosition lastPosition = timestamp;
var filterDefinition = CreateFilter(streamFilter, lastPosition);
var find = var find =
Collection.Find(filterDefinition, Batching.Options) Collection.Find(CreateFilter(filter, lastPosition), Batching.Options)
.Limit(take).Sort(Sort.Descending(x => x.Timestamp).Ascending(x => x.EventStream)); .Limit(take).Sort(Sort.Descending(x => x.Timestamp).Ascending(x => x.EventStream));
var taken = 0; var taken = 0;
@ -158,12 +128,12 @@ public partial class MongoEventStore : MongoRepositoryBase<MongoEventCommit>, IE
} }
} }
public async IAsyncEnumerable<StoredEvent> QueryAllAsync(string? streamFilter = null, string? position = null, int take = int.MaxValue, public async IAsyncEnumerable<StoredEvent> QueryAllAsync(StreamFilter filter, string? position = null, int take = int.MaxValue,
[EnumeratorCancellation] CancellationToken ct = default) [EnumeratorCancellation] CancellationToken ct = default)
{ {
StreamPosition lastPosition = position; StreamPosition lastPosition = position;
var filterDefinition = CreateFilter(streamFilter, lastPosition); var filterDefinition = CreateFilter(filter, lastPosition);
var find = var find =
Collection.Find(filterDefinition).SortBy(x => x.Timestamp).ThenByDescending(x => x.EventStream) Collection.Find(filterDefinition).SortBy(x => x.Timestamp).ThenByDescending(x => x.EventStream)
@ -197,15 +167,13 @@ public partial class MongoEventStore : MongoRepositoryBase<MongoEventCommit>, IE
return commits.OrderBy(x => x.EventStreamOffset).ThenBy(x => x.Timestamp).SelectMany(x => x.Filtered(streamPosition)).ToList(); return commits.OrderBy(x => x.EventStreamOffset).ThenBy(x => x.Timestamp).SelectMany(x => x.Filtered(streamPosition)).ToList();
} }
private static FilterDefinition<MongoEventCommit> CreateFilter(string? streamFilter, StreamPosition streamPosition) private static FilterDefinition<MongoEventCommit> CreateFilter(StreamFilter filter, StreamPosition streamPosition)
{ {
var filter = FilterExtensions.ByPosition(streamPosition); return Filter.And(FilterExtensions.ByPosition(streamPosition), FilterExtensions.ByStream(filter));
}
if (streamFilter != null)
{
return Filter.And(filter, FilterExtensions.ByStream(streamFilter));
}
return filter; private static FilterDefinition<MongoEventCommit> CreateFilter(StreamFilter filter, long streamPosition)
{
return Filter.And(FilterExtensions.ByStream(filter), FilterExtensions.ByOffset(streamPosition));
} }
} }

14
backend/src/Squidex.Infrastructure.MongoDb/EventSourcing/MongoEventStore_Writer.cs

@ -17,20 +17,12 @@ public partial class MongoEventStore
private const int MaxWriteAttempts = 20; private const int MaxWriteAttempts = 20;
private static readonly BsonTimestamp EmptyTimestamp = new BsonTimestamp(0); private static readonly BsonTimestamp EmptyTimestamp = new BsonTimestamp(0);
public Task DeleteStreamAsync(string streamName, public Task DeleteAsync(StreamFilter filter,
CancellationToken ct = default) CancellationToken ct = default)
{ {
Guard.NotNullOrEmpty(streamName); Guard.NotDefault(filter);
return Collection.DeleteManyAsync(x => x.EventStream == streamName, ct);
}
public Task DeleteAsync(string streamFilter,
CancellationToken ct = default)
{
Guard.NotNullOrEmpty(streamFilter);
return Collection.DeleteManyAsync(FilterExtensions.ByStream(streamFilter), ct); return Collection.DeleteManyAsync(FilterExtensions.ByStream(filter), ct);
} }
public async Task AppendAsync(Guid commitId, string streamName, long expectedVersion, ICollection<EventData> events, public async Task AppendAsync(Guid commitId, string streamName, long expectedVersion, ICollection<EventData> events,

4
backend/src/Squidex.Infrastructure/Commands/Rebuilder.cs

@ -49,14 +49,14 @@ public class Rebuilder
return domainObject; return domainObject;
} }
public virtual Task RebuildAsync<T, TState>(string filter, int batchSize, public virtual Task RebuildAsync<T, TState>(StreamFilter filter, int batchSize,
CancellationToken ct = default) CancellationToken ct = default)
where T : DomainObject<TState> where TState : class, IDomainState<TState>, new() where T : DomainObject<TState> where TState : class, IDomainState<TState>, new()
{ {
return RebuildAsync<T, TState>(filter, batchSize, 0, ct); return RebuildAsync<T, TState>(filter, batchSize, 0, ct);
} }
public virtual async Task RebuildAsync<T, TState>(string filter, int batchSize, double errorThreshold, public virtual async Task RebuildAsync<T, TState>(StreamFilter filter, int batchSize, double errorThreshold,
CancellationToken ct = default) CancellationToken ct = default)
where T : DomainObject<TState> where TState : class, IDomainState<TState>, new() where T : DomainObject<TState> where TState : class, IDomainState<TState>, new()
{ {

4
backend/src/Squidex.Infrastructure/EventSourcing/IEventConsumer.cs

@ -13,9 +13,9 @@ public interface IEventConsumer
int BatchSize => 1; int BatchSize => 1;
string Name { get; } string Name => GetType().Name;
string EventsFilter => ".*"; StreamFilter EventsFilter => default;
bool StartLatest => false; bool StartLatest => false;

28
backend/src/Squidex.Infrastructure/EventSourcing/IEventStore.cs

@ -11,28 +11,25 @@ namespace Squidex.Infrastructure.EventSourcing;
public interface IEventStore public interface IEventStore
{ {
Task<IReadOnlyList<StoredEvent>> QueryReverseAsync(string streamName, int take = int.MaxValue, Task<IReadOnlyList<StoredEvent>> QueryStreamReverseAsync(string streamName, int take = int.MaxValue,
CancellationToken ct = default); CancellationToken ct = default);
Task<IReadOnlyList<StoredEvent>> QueryAsync(string streamName, long afterStreamPosition = EtagVersion.Empty, Task<IReadOnlyList<StoredEvent>> QueryStreamAsync(string streamName, long afterStreamPosition = EtagVersion.Empty,
CancellationToken ct = default); CancellationToken ct = default);
IAsyncEnumerable<StoredEvent> QueryAllReverseAsync(string? streamFilter = null, Instant timestamp = default, int take = int.MaxValue, IAsyncEnumerable<StoredEvent> QueryAllReverseAsync(StreamFilter filter, Instant timestamp = default, int take = int.MaxValue,
CancellationToken ct = default); CancellationToken ct = default);
IAsyncEnumerable<StoredEvent> QueryAllAsync(string? streamFilter = null, string? position = null, int take = int.MaxValue, IAsyncEnumerable<StoredEvent> QueryAllAsync(StreamFilter filter, string? position = null, int take = int.MaxValue,
CancellationToken ct = default); CancellationToken ct = default);
Task AppendAsync(Guid commitId, string streamName, long expectedVersion, ICollection<EventData> events, Task AppendAsync(Guid commitId, string streamName, long expectedVersion, ICollection<EventData> events,
CancellationToken ct = default); CancellationToken ct = default);
Task DeleteAsync(string streamFilter, Task DeleteAsync(StreamFilter filter,
CancellationToken ct = default); CancellationToken ct = default);
Task DeleteStreamAsync(string streamName, IEventSubscription CreateSubscription(IEventSubscriber<StoredEvent> eventSubscriber, StreamFilter filter, string? position = null);
CancellationToken ct = default);
IEventSubscription CreateSubscription(IEventSubscriber<StoredEvent> eventSubscriber, string? streamFilter = null, string? position = null);
async Task AppendUnsafeAsync(IEnumerable<EventCommit> commits, async Task AppendUnsafeAsync(IEnumerable<EventCommit> commits,
CancellationToken ct = default) CancellationToken ct = default)
@ -42,17 +39,4 @@ public interface IEventStore
await AppendAsync(commit.Id, commit.StreamName, commit.Offset, commit.Events, ct); await AppendAsync(commit.Id, commit.StreamName, commit.Offset, commit.Events, ct);
} }
} }
async Task<IReadOnlyDictionary<string, IReadOnlyList<StoredEvent>>> QueryManyAsync(IEnumerable<string> streamNames,
CancellationToken ct = default)
{
var result = new Dictionary<string, IReadOnlyList<StoredEvent>>();
foreach (var streamName in streamNames)
{
result[streamName] = await QueryAsync(streamName, EtagVersion.Empty, ct);
}
return result;
}
} }

2
backend/src/Squidex.Infrastructure/EventSourcing/PollingSubscription.cs

@ -17,7 +17,7 @@ public sealed class PollingSubscription : IEventSubscription
public PollingSubscription( public PollingSubscription(
IEventStore eventStore, IEventStore eventStore,
IEventSubscriber<StoredEvent> eventSubscriber, IEventSubscriber<StoredEvent> eventSubscriber,
string? streamFilter, StreamFilter streamFilter,
string? position) string? position)
{ {
timer = new CompletionTimer(5000, async ct => timer = new CompletionTimer(5000, async ct =>

35
backend/src/Squidex.Infrastructure/EventSourcing/StreamFilter.cs

@ -5,17 +5,38 @@
// All rights reserved. Licensed under the MIT license. // All rights reserved. Licensed under the MIT license.
// ========================================================================== // ==========================================================================
using System.Diagnostics.CodeAnalysis; using Squidex.Infrastructure.Collections;
namespace Squidex.Infrastructure.EventSourcing; namespace Squidex.Infrastructure.EventSourcing;
public static class StreamFilter public readonly record struct StreamFilter
{ {
public static bool IsAll([NotNullWhen(false)] string? filter) public ReadonlyList<string>? Prefixes { get; }
public StreamFilterKind Kind { get; }
public StreamFilter(StreamFilterKind kind, params string[] prefixes)
{
Kind = kind;
if (prefixes.Length > 0)
{
Prefixes = prefixes.ToReadonlyList();
}
}
public static StreamFilter Prefix(params string[] prefixes)
{
return new StreamFilter(StreamFilterKind.MatchStart, prefixes);
}
public static StreamFilter Name(params string[] prefixes)
{
return new StreamFilter(StreamFilterKind.MatchFull, prefixes);
}
public static StreamFilter All()
{ {
return string.IsNullOrWhiteSpace(filter) return default;
|| string.Equals(filter, ".*", StringComparison.OrdinalIgnoreCase)
|| string.Equals(filter, "(.*)", StringComparison.OrdinalIgnoreCase)
|| string.Equals(filter, "(.*?)", StringComparison.OrdinalIgnoreCase);
} }
} }

14
backend/src/Squidex.Infrastructure/EventSourcing/StreamFilterKind.cs

@ -0,0 +1,14 @@
// ==========================================================================
// Squidex Headless CMS
// ==========================================================================
// Copyright (c) Squidex UG (haftungsbeschraenkt)
// All rights reserved. Licensed under the MIT license.
// ==========================================================================
namespace Squidex.Infrastructure.EventSourcing;
public enum StreamFilterKind
{
MatchFull,
MatchStart
}

15
backend/src/Squidex.Infrastructure/States/BatchContext.cs

@ -50,22 +50,27 @@ public sealed class BatchContext<T> : IBatchContext<T>
public async Task LoadAsync(IEnumerable<DomainId> ids) public async Task LoadAsync(IEnumerable<DomainId> ids)
{ {
var streamNames = ids.ToDictionary(x => x, x => eventStreamNames.GetStreamName(owner, x.ToString())); var streamNames = ids.ToDictionary(
x => x,
x => eventStreamNames.GetStreamName(owner, x.ToString()));
if (streamNames.Count == 0) if (streamNames.Count == 0)
{ {
return; return;
} }
var streams = await eventStore.QueryManyAsync(streamNames.Values); var eventsResults = await eventStore.QueryAllAsync(StreamFilter.Name(streamNames.Values.ToArray())).ToListAsync();
var eventsByStream = eventsResults.ToLookup(x => x.StreamName);
foreach (var (id, streamName) in streamNames) foreach (var (id, streamName) in streamNames)
{ {
if (streams.TryGetValue(streamName, out var data)) var byStream = eventsByStream[streamName].ToList();
if (byStream.Count > 0)
{ {
var stream = data.Select(eventFormatter.ParseIfKnown).NotNull().ToList(); var parsed = byStream.Select(eventFormatter.ParseIfKnown).NotNull().ToList();
events[id] = (data.Count - 1, stream); events[id] = (byStream.Count - 1, parsed);
} }
else else
{ {

4
backend/src/Squidex.Infrastructure/States/Persistence.cs

@ -82,7 +82,7 @@ internal sealed class Persistence<T> : IPersistence<T>
{ {
using (Telemetry.Activities.StartActivity("Persistence/ReadEvents")) using (Telemetry.Activities.StartActivity("Persistence/ReadEvents"))
{ {
await eventStore.DeleteStreamAsync(streamName.Value, ct); await eventStore.DeleteAsync(StreamFilter.Name(streamName.Value), ct);
} }
versionEvents = EtagVersion.Empty; versionEvents = EtagVersion.Empty;
@ -144,7 +144,7 @@ internal sealed class Persistence<T> : IPersistence<T>
private async Task ReadEventsAsync( private async Task ReadEventsAsync(
CancellationToken ct) CancellationToken ct)
{ {
var events = await eventStore.QueryAsync(streamName.Value, versionEvents, ct); var events = await eventStore.QueryStreamAsync(streamName.Value, versionEvents, ct);
var isStopped = false; var isStopped = false;

2
backend/tests/Squidex.Domain.Apps.Core.Tests/Operations/Subscriptions/SubscriptionPublisherTests.cs

@ -29,7 +29,7 @@ public class SubscriptionPublisherTests
[Fact] [Fact]
public void Should_return_content_and_asset_filter_for_events_filter() public void Should_return_content_and_asset_filter_for_events_filter()
{ {
Assert.Equal("^(content-|asset-)", sut.EventsFilter); Assert.Equal(StreamFilter.Prefix("content-", "asset-"), sut.EventsFilter);
} }
[Fact] [Fact]

4
backend/tests/Squidex.Domain.Apps.Entities.Tests/Apps/AppEventDeleterTests.cs

@ -33,7 +33,9 @@ public class AppEventDeleterTests : GivenContext
{ {
await sut.DeleteAppAsync(App, CancellationToken); await sut.DeleteAppAsync(App, CancellationToken);
A.CallTo(() => eventStore.DeleteAsync($"^[a-zA-Z0-9]-{AppId.Id}", A<CancellationToken>._)) var streamFilter = StreamFilter.Prefix($"[a-zA-Z0-9]-{AppId.Id}");
A.CallTo(() => eventStore.DeleteAsync(streamFilter, A<CancellationToken>._))
.MustNotHaveHappened(); .MustNotHaveHappened();
} }
} }

4
backend/tests/Squidex.Domain.Apps.Entities.Tests/Apps/AppPermanentDeleterTests.cs

@ -29,7 +29,7 @@ public class AppPermanentDeleterTests : GivenContext
[Fact] [Fact]
public void Should_return_assets_filter_for_events_filter() public void Should_return_assets_filter_for_events_filter()
{ {
Assert.Equal("^app-", sut.EventsFilter); Assert.Equal(StreamFilter.Prefix("app-"), sut.EventsFilter);
} }
[Fact] [Fact]
@ -41,7 +41,7 @@ public class AppPermanentDeleterTests : GivenContext
[Fact] [Fact]
public void Should_return_type_name_for_name() public void Should_return_type_name_for_name()
{ {
Assert.Equal(nameof(AppPermanentDeleter), sut.Name); Assert.Equal(nameof(AppPermanentDeleter), ((IEventConsumer)sut).Name);
} }
[Fact] [Fact]

4
backend/tests/Squidex.Domain.Apps.Entities.Tests/Assets/AssetPermanentDeleterTests.cs

@ -27,7 +27,7 @@ public class AssetPermanentDeleterTests : GivenContext
[Fact] [Fact]
public void Should_return_assets_filter_for_events_filter() public void Should_return_assets_filter_for_events_filter()
{ {
Assert.Equal("^asset-", sut.EventsFilter); Assert.Equal(StreamFilter.Prefix("asset-"), sut.EventsFilter);
} }
[Fact] [Fact]
@ -39,7 +39,7 @@ public class AssetPermanentDeleterTests : GivenContext
[Fact] [Fact]
public void Should_return_type_name_for_name() public void Should_return_type_name_for_name()
{ {
Assert.Equal(nameof(AssetPermanentDeleter), sut.Name); Assert.Equal(nameof(AssetPermanentDeleter), ((IEventConsumer)sut).Name);
} }
[Fact] [Fact]

4
backend/tests/Squidex.Domain.Apps.Entities.Tests/Assets/AssetUsageTrackerTests.cs

@ -35,7 +35,7 @@ public class AssetUsageTrackerTests : GivenContext
[Fact] [Fact]
public void Should_return_assets_filter_for_events_filter() public void Should_return_assets_filter_for_events_filter()
{ {
Assert.Equal("^asset-", sut.EventsFilter); Assert.Equal(StreamFilter.Prefix("asset-"), sut.EventsFilter);
} }
[Fact] [Fact]
@ -47,7 +47,7 @@ public class AssetUsageTrackerTests : GivenContext
[Fact] [Fact]
public void Should_return_type_name_for_name() public void Should_return_type_name_for_name()
{ {
Assert.Equal(nameof(AssetUsageTracker), sut.Name); Assert.Equal(nameof(AssetUsageTracker), ((IEventConsumer)sut).Name);
} }
[Fact] [Fact]

4
backend/tests/Squidex.Domain.Apps.Entities.Tests/Assets/RecursiveDeleterTests.cs

@ -33,7 +33,7 @@ public class RecursiveDeleterTests : GivenContext
[Fact] [Fact]
public void Should_return_assets_filter_for_events_filter() public void Should_return_assets_filter_for_events_filter()
{ {
Assert.Equal("^assetFolder-", sut.EventsFilter); Assert.Equal(StreamFilter.Prefix("assetFolder-"), sut.EventsFilter);
} }
[Fact] [Fact]
@ -45,7 +45,7 @@ public class RecursiveDeleterTests : GivenContext
[Fact] [Fact]
public void Should_return_type_name_for_name() public void Should_return_type_name_for_name()
{ {
Assert.Equal(nameof(RecursiveDeleter), sut.Name); Assert.Equal(nameof(RecursiveDeleter), ((IEventConsumer)sut).Name);
} }
[Fact] [Fact]

4
backend/tests/Squidex.Domain.Apps.Entities.Tests/Assets/RepairFilesTests.cs

@ -122,7 +122,9 @@ public class RepairFilesTests : GivenContext
.Returns(null); .Returns(null);
} }
A.CallTo(() => eventStore.QueryAllAsync("^asset\\-", null, int.MaxValue, CancellationToken)) var streamFilter = StreamFilter.Prefix("asset-");
A.CallTo(() => eventStore.QueryAllAsync(streamFilter, null, int.MaxValue, CancellationToken))
.Returns(storedEvents.ToAsyncEnumerable()); .Returns(storedEvents.ToAsyncEnumerable());
} }
} }

6
backend/tests/Squidex.Domain.Apps.Entities.Tests/Contents/Text/TextIndexerTestsBase.cs

@ -49,6 +49,12 @@ public abstract class TextIndexerTestsBase : GivenContext
public abstract ITextIndex CreateIndex(); public abstract ITextIndex CreateIndex();
[Fact]
public void Should_return_content_filter_for_events_filter()
{
Assert.Equal(StreamFilter.Prefix("content-"), Sut.EventsFilter);
}
[Fact] [Fact]
public async Task Should_search_with_fuzzy() public async Task Should_search_with_fuzzy()
{ {

19
backend/tests/Squidex.Domain.Apps.Entities.Tests/Invitation/InvitationEventConsumerTests.cs

@ -9,6 +9,7 @@ using Microsoft.Extensions.Logging;
using NodaTime; using NodaTime;
using Squidex.Domain.Apps.Core.TestHelpers; using Squidex.Domain.Apps.Core.TestHelpers;
using Squidex.Domain.Apps.Entities.Apps; using Squidex.Domain.Apps.Entities.Apps;
using Squidex.Domain.Apps.Entities.Assets;
using Squidex.Domain.Apps.Entities.Notifications; using Squidex.Domain.Apps.Entities.Notifications;
using Squidex.Domain.Apps.Entities.Teams; using Squidex.Domain.Apps.Entities.Teams;
using Squidex.Domain.Apps.Entities.TestHelpers; using Squidex.Domain.Apps.Entities.TestHelpers;
@ -45,6 +46,24 @@ public class InvitationEventConsumerTests : GivenContext
sut = new InvitationEventConsumer(AppProvider, userNotifications, userResolver, log); sut = new InvitationEventConsumer(AppProvider, userNotifications, userResolver, log);
} }
[Fact]
public void Should_return_app_filter_for_events_filter()
{
Assert.Equal(StreamFilter.Prefix("app-"), sut.EventsFilter);
}
[Fact]
public async Task Should_do_nothing_on_clear()
{
await ((IEventConsumer)sut).ClearAsync();
}
[Fact]
public void Should_return_custom_name()
{
Assert.Equal("NotificationEmailSender", sut.Name);
}
[Fact] [Fact]
public async Task Should_not_send_app_email_if_contributors_assigned_by_clients() public async Task Should_not_send_app_email_if_contributors_assigned_by_clients()
{ {

2
backend/tests/Squidex.Domain.Apps.Entities.Tests/Rules/RuleEnqueuerTests.cs

@ -52,7 +52,7 @@ public class RuleEnqueuerTests : GivenContext
[Fact] [Fact]
public void Should_return_wildcard_filter_for_events_filter() public void Should_return_wildcard_filter_for_events_filter()
{ {
Assert.Equal(".*", ((IEventConsumer)sut).EventsFilter); Assert.Equal(default, ((IEventConsumer)sut).EventsFilter);
} }
[Fact] [Fact]

22
backend/tests/Squidex.Infrastructure.Tests/EventSourcing/Consume/EventConsumerProcessorTests.cs

@ -62,7 +62,7 @@ public class EventConsumerProcessorTests
} }
}; };
A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<string>._, A<string>._)) A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<StreamFilter>._, A<string>._))
.Returns(eventSubscription); .Returns(eventSubscription);
A.CallTo(() => eventConsumer.Name) A.CallTo(() => eventConsumer.Name)
@ -101,15 +101,17 @@ public class EventConsumerProcessorTests
{ {
state.Snapshot = new EventConsumerState(); state.Snapshot = new EventConsumerState();
var filter = StreamFilter.Name("my-filter");
A.CallTo(() => eventConsumer.StartLatest) A.CallTo(() => eventConsumer.StartLatest)
.Returns(true); .Returns(true);
A.CallTo(() => eventConsumer.EventsFilter) A.CallTo(() => eventConsumer.EventsFilter)
.Returns("my-filter"); .Returns(filter);
var latestPosition = "LATEST"; var latestPosition = "LATEST";
A.CallTo(() => eventStore.QueryAllReverseAsync("my-filter", default, 1, A<CancellationToken>._)) A.CallTo(() => eventStore.QueryAllReverseAsync(filter, default, 1, A<CancellationToken>._))
.Returns(Enumerable.Repeat(new StoredEvent("Stream", latestPosition, 1, eventData), 1).ToAsyncEnumerable()); .Returns(Enumerable.Repeat(new StoredEvent("Stream", latestPosition, 1, eventData), 1).ToAsyncEnumerable());
await sut.InitializeAsync(default); await sut.InitializeAsync(default);
@ -129,7 +131,7 @@ public class EventConsumerProcessorTests
AssertGrainState(isStopped: true, position: initialPosition); AssertGrainState(isStopped: true, position: initialPosition);
A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<string>._, A<string>._)) A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<StreamFilter>._, A<string>._))
.MustNotHaveHappened(); .MustNotHaveHappened();
} }
@ -143,7 +145,7 @@ public class EventConsumerProcessorTests
AssertGrainState(isStopped: false, position: initialPosition); AssertGrainState(isStopped: false, position: initialPosition);
A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<string>._, A<string>._)) A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<StreamFilter>._, A<string>._))
.MustHaveHappenedOnceExactly(); .MustHaveHappenedOnceExactly();
} }
@ -159,7 +161,7 @@ public class EventConsumerProcessorTests
AssertGrainState(isStopped: false, position: initialPosition); AssertGrainState(isStopped: false, position: initialPosition);
A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<string>._, A<string>._)) A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<StreamFilter>._, A<string>._))
.MustHaveHappenedOnceExactly(); .MustHaveHappenedOnceExactly();
} }
@ -173,7 +175,7 @@ public class EventConsumerProcessorTests
AssertGrainState(isStopped: false, position: initialPosition); AssertGrainState(isStopped: false, position: initialPosition);
A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<string>._, A<string>._)) A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<StreamFilter>._, A<string>._))
.MustHaveHappenedOnceExactly(); .MustHaveHappenedOnceExactly();
} }
@ -219,10 +221,10 @@ public class EventConsumerProcessorTests
A.CallTo(() => eventSubscription.Dispose()) A.CallTo(() => eventSubscription.Dispose())
.MustHaveHappenedOnceExactly(); .MustHaveHappenedOnceExactly();
A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<string>._, state.Snapshot.Position)) A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<StreamFilter>._, state.Snapshot.Position))
.MustHaveHappenedOnceExactly(); .MustHaveHappenedOnceExactly();
A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<string>._, null)) A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<StreamFilter>._, null))
.MustHaveHappenedOnceExactly(); .MustHaveHappenedOnceExactly();
} }
@ -530,7 +532,7 @@ public class EventConsumerProcessorTests
A.CallTo(() => eventSubscription.Dispose()) A.CallTo(() => eventSubscription.Dispose())
.MustHaveHappenedOnceExactly(); .MustHaveHappenedOnceExactly();
A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<string>._, A<string>._)) A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<StreamFilter>._, A<string>._))
.MustHaveHappened(2, Times.Exactly); .MustHaveHappened(2, Times.Exactly);
} }

62
backend/tests/Squidex.Infrastructure.Tests/EventSourcing/EventStoreTests.cs

@ -7,6 +7,7 @@
using System.Globalization; using System.Globalization;
using System.Text.RegularExpressions; using System.Text.RegularExpressions;
using EventStore.Client;
namespace Squidex.Infrastructure.EventSourcing; namespace Squidex.Infrastructure.EventSourcing;
@ -90,6 +91,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
public async Task Should_append_events() public async Task Should_append_events()
{ {
var streamName = $"test-{Guid.NewGuid()}"; var streamName = $"test-{Guid.NewGuid()}";
var streamFilter = StreamFilter.Name(streamName);
var commit1 = new[] var commit1 = new[]
{ {
@ -107,7 +109,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
await Sut.AppendAsync(Guid.NewGuid(), streamName, EtagVersion.Any, commit2); await Sut.AppendAsync(Guid.NewGuid(), streamName, EtagVersion.Any, commit2);
var readEvents1 = await QueryAsync(streamName); var readEvents1 = await QueryAsync(streamName);
var readEvents2 = await QueryAllAsync(streamName); var readEvents2 = await QueryAllAsync(streamFilter);
var expected = new[] var expected = new[]
{ {
@ -125,6 +127,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
public async Task Should_append_events_unsafe() public async Task Should_append_events_unsafe()
{ {
var streamName = $"test-{Guid.NewGuid()}"; var streamName = $"test-{Guid.NewGuid()}";
var streamFilter = StreamFilter.Name(streamName);
var commit1 = new[] var commit1 = new[]
{ {
@ -138,7 +141,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
}); });
var readEvents1 = await QueryAsync(streamName); var readEvents1 = await QueryAsync(streamName);
var readEvents2 = await QueryAllAsync(streamName); var readEvents2 = await QueryAllAsync(streamFilter);
var expected = new[] var expected = new[]
{ {
@ -154,6 +157,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
public async Task Should_subscribe_to_events() public async Task Should_subscribe_to_events()
{ {
var streamName = $"test-{Guid.NewGuid()}"; var streamName = $"test-{Guid.NewGuid()}";
var streamFilter = StreamFilter.Name(streamName);
var commit1 = new[] var commit1 = new[]
{ {
@ -161,7 +165,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
CreateEventData(2) CreateEventData(2)
}; };
var readEvents = await QueryWithSubscriptionAsync(streamName, async () => var readEvents = await QueryWithSubscriptionAsync(streamFilter, async () =>
{ {
await Sut.AppendAsync(Guid.NewGuid(), streamName, EtagVersion.Any, commit1); await Sut.AppendAsync(Guid.NewGuid(), streamName, EtagVersion.Any, commit1);
}); });
@ -179,6 +183,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
public async Task Should_subscribe_to_next_events() public async Task Should_subscribe_to_next_events()
{ {
var streamName = $"test-{Guid.NewGuid()}"; var streamName = $"test-{Guid.NewGuid()}";
var streamFilter = StreamFilter.Name(streamName);
var commit1 = new[] var commit1 = new[]
{ {
@ -187,7 +192,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
}; };
// Append and read in parallel. // Append and read in parallel.
await QueryWithSubscriptionAsync(streamName, async () => await QueryWithSubscriptionAsync(streamFilter, async () =>
{ {
await Sut.AppendAsync(Guid.NewGuid(), streamName, EtagVersion.Any, commit1); await Sut.AppendAsync(Guid.NewGuid(), streamName, EtagVersion.Any, commit1);
}); });
@ -199,7 +204,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
}; };
// Append and read in parallel. // Append and read in parallel.
var readEventsFromPosition = await QueryWithSubscriptionAsync(streamName, async () => var readEventsFromPosition = await QueryWithSubscriptionAsync(streamFilter, async () =>
{ {
await Sut.AppendAsync(Guid.NewGuid(), streamName, EtagVersion.Any, commit2); await Sut.AppendAsync(Guid.NewGuid(), streamName, EtagVersion.Any, commit2);
}); });
@ -210,7 +215,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
new StoredEvent(streamName, "Position", 3, commit2[1]) new StoredEvent(streamName, "Position", 3, commit2[1])
}; };
var readEventsFromBeginning = await QueryWithSubscriptionAsync(streamName, fromBeginning: true); var readEventsFromBeginning = await QueryWithSubscriptionAsync(streamFilter, fromBeginning: true);
var expectedFromBeginning = new[] var expectedFromBeginning = new[]
{ {
@ -228,12 +233,13 @@ public abstract class EventStoreTests<T> where T : IEventStore
public async Task Should_subscribe_with_parallel_writes() public async Task Should_subscribe_with_parallel_writes()
{ {
var streamName = $"test-{Guid.NewGuid()}"; var streamName = $"test-{Guid.NewGuid()}";
var streamFilter = StreamFilter.Prefix(streamName);
var numTasks = 50; var numTasks = 50;
var numEvents = 100; var numEvents = 100;
// Append and read in parallel. // Append and read in parallel.
var readEvents = await QueryWithSubscriptionAsync($"^{streamName}", async () => var readEvents = await QueryWithSubscriptionAsync(streamFilter, async () =>
{ {
await Parallel.ForEachAsync(Enumerable.Range(0, numTasks), async (i, ct) => await Parallel.ForEachAsync(Enumerable.Range(0, numTasks), async (i, ct) =>
{ {
@ -277,7 +283,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
await Sut.AppendAsync(Guid.NewGuid(), streamName1, EtagVersion.Any, stream1Commit); await Sut.AppendAsync(Guid.NewGuid(), streamName1, EtagVersion.Any, stream1Commit);
await Sut.AppendAsync(Guid.NewGuid(), streamName2, EtagVersion.Any, stream2Commit); await Sut.AppendAsync(Guid.NewGuid(), streamName2, EtagVersion.Any, stream2Commit);
var readEvents = await Sut.QueryManyAsync(new[] { streamName1, streamName2 }); var readEvents = await Sut.QueryAllAsync(StreamFilter.Name(streamName1, streamName2)).ToListAsync();
var expected1 = new[] var expected1 = new[]
{ {
@ -291,8 +297,8 @@ public abstract class EventStoreTests<T> where T : IEventStore
new StoredEvent(streamName2, "Position", 1, stream2Commit[1]) new StoredEvent(streamName2, "Position", 1, stream2Commit[1])
}; };
ShouldBeEquivalentTo(readEvents[streamName1], expected1); ShouldBeEquivalentTo(readEvents.Where(x => x.StreamName == streamName1), expected1);
ShouldBeEquivalentTo(readEvents[streamName2], expected2); ShouldBeEquivalentTo(readEvents.Where(x => x.StreamName == streamName2), expected2);
} }
[Theory] [Theory]
@ -300,7 +306,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
[InlineData(5, 30)] [InlineData(5, 30)]
[InlineData(5, 300)] [InlineData(5, 300)]
[InlineData(5, 3000)] [InlineData(5, 3000)]
public async Task Should_read_events_from_offset(int commits, int count) public async Task Should_query_events_from_offset(int commits, int count)
{ {
var streamName = $"test-{Guid.NewGuid()}"; var streamName = $"test-{Guid.NewGuid()}";
@ -308,7 +314,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
var readEvents0 = await QueryAsync(streamName); var readEvents0 = await QueryAsync(streamName);
var readEvents1 = await QueryAsync(streamName, count - 2); var readEvents1 = await QueryAsync(streamName, count - 2);
var readEvents2 = await QueryAllAsync(streamName, readEvents0[^2].EventPosition); var readEvents2 = await QueryAllAsync(default, readEvents0[^2].EventPosition);
var expected = new[] var expected = new[]
{ {
@ -322,7 +328,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
[Theory] [Theory]
[InlineData(5, 30)] [InlineData(5, 30)]
[InlineData(5, 300)] [InlineData(5, 300)]
public async Task Should_read_reverse(int commits, int count) public async Task Should_query_reverse(int commits, int count)
{ {
var streamName = $"test-{Guid.NewGuid()}"; var streamName = $"test-{Guid.NewGuid()}";
@ -332,7 +338,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
for (var take = 0; take < count; take += count / 10) for (var take = 0; take < count; take += count / 10)
{ {
var eventsExpected = eventsStored.TakeLast(take).ToArray(); var eventsExpected = eventsStored.TakeLast(take).ToArray();
var eventsQueried = await Sut.QueryReverseAsync(streamName, take); var eventsQueried = await Sut.QueryStreamReverseAsync(streamName, take);
ShouldBeEquivalentTo(eventsQueried, eventsExpected); ShouldBeEquivalentTo(eventsQueried, eventsExpected);
} }
@ -342,10 +348,10 @@ public abstract class EventStoreTests<T> where T : IEventStore
[InlineData(5, 30)] [InlineData(5, 30)]
[InlineData(5, 300)] [InlineData(5, 300)]
[InlineData(5, 3000)] [InlineData(5, 3000)]
public async Task Should_read_all_reverse_by_name(int commits, int count) public async Task Should_query_all_reverse_by_names(int commits, int count)
{ {
var streamName = $"test-{Guid.NewGuid()}"; var streamName = $"test-{Guid.NewGuid()}";
var streamFilter = streamName; var streamFilter = StreamFilter.Name(streamName, "invalid");
var eventsWritten = await AppendEventsAsync(streamName, count, commits); var eventsWritten = await AppendEventsAsync(streamName, count, commits);
var eventsStored = eventsWritten.Select((x, i) => new StoredEvent(streamName, "Position", i, x)).ToArray(); var eventsStored = eventsWritten.Select((x, i) => new StoredEvent(streamName, "Position", i, x)).ToArray();
@ -353,7 +359,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
for (var take = 0; take < count; take += count / 10) for (var take = 0; take < count; take += count / 10)
{ {
var eventsExpected = eventsStored.Reverse().Take(take).ToArray(); var eventsExpected = eventsStored.Reverse().Take(take).ToArray();
var eventsQueried = await Sut.QueryAllReverseAsync(streamName, default, take).ToArrayAsync(); var eventsQueried = await Sut.QueryAllReverseAsync(streamFilter, default, take).ToArrayAsync();
ShouldBeEquivalentTo(eventsQueried, eventsExpected); ShouldBeEquivalentTo(eventsQueried, eventsExpected);
} }
@ -363,10 +369,10 @@ public abstract class EventStoreTests<T> where T : IEventStore
[InlineData(5, 30)] [InlineData(5, 30)]
[InlineData(5, 300)] [InlineData(5, 300)]
[InlineData(5, 3000)] [InlineData(5, 3000)]
public async Task Should_read_all_reverse_by_filter(int commits, int count) public async Task Should_query_all_reverse_by_filter(int commits, int count)
{ {
var streamName = $"test-{Guid.NewGuid()}-suffix"; var streamName = $"test-{Guid.NewGuid()}-suffix";
var streamFilter = $"^{Regex.Escape(streamName)}"; var streamFilter = StreamFilter.Prefix(streamName[..^7], "invalid");
var eventsWritten = await AppendEventsAsync(streamName, count, commits); var eventsWritten = await AppendEventsAsync(streamName, count, commits);
var eventsStored = eventsWritten.Select((x, i) => new StoredEvent(streamName, "Position", i, x)).ToArray(); var eventsStored = eventsWritten.Select((x, i) => new StoredEvent(streamName, "Position", i, x)).ToArray();
@ -387,7 +393,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
public async Task Should_read_all_reverse(int commits, int count) public async Task Should_read_all_reverse(int commits, int count)
{ {
var streamName = $"test-{Guid.NewGuid()}-suffix"; var streamName = $"test-{Guid.NewGuid()}-suffix";
var streamFilter = ".*"; var streamFilter = default(StreamFilter);
await AppendEventsAsync(streamName, count, commits); await AppendEventsAsync(streamName, count, commits);
@ -403,6 +409,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
public async Task Should_delete_by_filter() public async Task Should_delete_by_filter()
{ {
var streamName = $"test-{Guid.NewGuid()}"; var streamName = $"test-{Guid.NewGuid()}";
var streamFilter = StreamFilter.Prefix($"{streamName[..10]}");
await AppendEventsAsync(streamName, 2, 1); await AppendEventsAsync(streamName, 2, 1);
@ -410,7 +417,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
for (var i = 0; i < 5; i++) for (var i = 0; i < 5; i++)
{ {
await Sut.DeleteAsync($"^{streamName[..10]}"); await Sut.DeleteAsync(streamFilter);
readEvents = await QueryAsync(streamName); readEvents = await QueryAsync(streamName);
@ -427,9 +434,10 @@ public abstract class EventStoreTests<T> where T : IEventStore
} }
[Fact] [Fact]
public async Task Should_delete_stream() public async Task Should_delete_by_name()
{ {
var streamName = $"test-{Guid.NewGuid()}"; var streamName = $"test-{Guid.NewGuid()}";
var streamFilter = StreamFilter.Name(streamName);
await AppendEventsAsync(streamName, 2, 1); await AppendEventsAsync(streamName, 2, 1);
@ -437,7 +445,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
for (var i = 0; i < 5; i++) for (var i = 0; i < 5; i++)
{ {
await Sut.DeleteStreamAsync(streamName); await Sut.DeleteAsync(streamFilter);
readEvents = await QueryAsync(streamName); readEvents = await QueryAsync(streamName);
@ -455,12 +463,12 @@ public abstract class EventStoreTests<T> where T : IEventStore
private async Task<IReadOnlyList<StoredEvent>> QueryAsync(string streamName, long position = EtagVersion.Any) private async Task<IReadOnlyList<StoredEvent>> QueryAsync(string streamName, long position = EtagVersion.Any)
{ {
return await Sut.QueryAsync(streamName, position); return await Sut.QueryStreamAsync(streamName, position);
} }
private async Task<IReadOnlyList<StoredEvent>?> QueryAllAsync(string? streamFilter = null, string? position = null) private async Task<IReadOnlyList<StoredEvent>?> QueryAllAsync(StreamFilter filter, string? position = null)
{ {
return await Sut.QueryAllAsync(streamFilter, position).ToListAsync(); return await Sut.QueryAllAsync(filter, position).ToListAsync();
} }
private static EventData CreateEventData(int i) private static EventData CreateEventData(int i)
@ -473,7 +481,7 @@ public abstract class EventStoreTests<T> where T : IEventStore
return new EventData($"Type{i}", headers, i.ToString(CultureInfo.InvariantCulture)); return new EventData($"Type{i}", headers, i.ToString(CultureInfo.InvariantCulture));
} }
private async Task<IReadOnlyList<StoredEvent>?> QueryWithSubscriptionAsync(string streamFilter, private async Task<IReadOnlyList<StoredEvent>?> QueryWithSubscriptionAsync(StreamFilter streamFilter,
Func<Task>? subscriptionRunning = null, bool fromBeginning = false) Func<Task>? subscriptionRunning = null, bool fromBeginning = false)
{ {
var subscriber = new EventSubscriber(); var subscriber = new EventSubscriber();

5
backend/tests/Squidex.Infrastructure.Tests/EventSourcing/MongoEventStoreTests.cs

@ -41,9 +41,8 @@ public abstract class MongoEventStoreTests : EventStoreTests<MongoEventStore>, I
Assert.All(queries, query => Assert.All(queries, query =>
{ {
Assert.Equal(query.NumDocuments, query.DocsExamined); Assert.InRange(query.DocsExamined, 0, query.NumDocuments);
Assert.True(query.KeysExamined >= query.NumDocuments); Assert.InRange(query.KeysExamined, query.NumDocuments, (Math.Max(1, query.NumDocuments) * 2) + 1);
Assert.True(query.KeysExamined <= query.NumDocuments * 2);
}); });
} }
} }

2
backend/tests/Squidex.Infrastructure.Tests/EventSourcing/MongoParallelInsertTests.cs

@ -38,7 +38,7 @@ public sealed class MongoParallelInsertTests : IClassFixture<MongoEventStoreFixt
public string Name { get; } = RandomHash.Simple(); public string Name { get; } = RandomHash.Simple();
public string EventsFilter => $"^{Name}"; public StreamFilter EventsFilter => StreamFilter.Prefix(Name);
public Task Completed => tcs.Task; public Task Completed => tcs.Task;

2
backend/tests/Squidex.Infrastructure.Tests/EventSourcing/PollingSubscriptionTests.cs

@ -13,8 +13,8 @@ public class PollingSubscriptionTests
{ {
private readonly IEventStore eventStore = A.Fake<IEventStore>(); private readonly IEventStore eventStore = A.Fake<IEventStore>();
private readonly IEventSubscriber<StoredEvent> eventSubscriber = A.Fake<IEventSubscriber<StoredEvent>>(); private readonly IEventSubscriber<StoredEvent> eventSubscriber = A.Fake<IEventSubscriber<StoredEvent>>();
private readonly StreamFilter filter = StreamFilter.Name("my-stream");
private readonly string position = Guid.NewGuid().ToString(); private readonly string position = Guid.NewGuid().ToString();
private readonly string filter = "^my-stream";
[Fact] [Fact]
public async Task Should_subscribe_on_start() public async Task Should_subscribe_on_start()

8
backend/tests/Squidex.Infrastructure.Tests/EventSourcing/RetrySubscriptionTests.cs

@ -17,10 +17,10 @@ public class RetrySubscriptionTests
public RetrySubscriptionTests() public RetrySubscriptionTests()
{ {
A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<string>._, A<string>._)) A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<StreamFilter>._, A<string>._))
.Returns(eventSubscription); .Returns(eventSubscription);
sut = new RetrySubscription<StoredEvent>(eventSubscriber, s => eventStore.CreateSubscription(s)) { ReconnectWaitMs = 50 }; sut = new RetrySubscription<StoredEvent>(eventSubscriber, s => eventStore.CreateSubscription(s, default)) { ReconnectWaitMs = 50 };
sutSubscriber = sut; sutSubscriber = sut;
} }
@ -29,7 +29,7 @@ public class RetrySubscriptionTests
{ {
sut.Dispose(); sut.Dispose();
A.CallTo(() => eventStore.CreateSubscription(sut, A<string>._, A<string>._)) A.CallTo(() => eventStore.CreateSubscription(sut, A<StreamFilter>._, A<string>._))
.MustHaveHappened(); .MustHaveHappened();
} }
@ -47,7 +47,7 @@ public class RetrySubscriptionTests
A.CallTo(() => eventSubscription.Dispose()) A.CallTo(() => eventSubscription.Dispose())
.MustHaveHappened(2, Times.Exactly); .MustHaveHappened(2, Times.Exactly);
A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<string>._, A<string>._)) A.CallTo(() => eventStore.CreateSubscription(A<IEventSubscriber<StoredEvent>>._, A<StreamFilter>._, A<string>._))
.MustHaveHappened(2, Times.Exactly); .MustHaveHappened(2, Times.Exactly);
A.CallTo(() => eventSubscriber.OnErrorAsync(eventSubscription, A<Exception>._)) A.CallTo(() => eventSubscriber.OnErrorAsync(eventSubscription, A<Exception>._))

27
backend/tests/Squidex.Infrastructure.Tests/EventSourcing/StreamFilterTests.cs

@ -0,0 +1,27 @@
// ==========================================================================
// Squidex Headless CMS
// ==========================================================================
// Copyright (c) Squidex UG (haftungsbeschraenkt)
// All rights reserved. Licensed under the MIT license.
// ==========================================================================
namespace Squidex.Infrastructure.EventSourcing;
public class StreamFilterTests
{
[Fact]
public void Should_simplify_input_to_default_filter()
{
var sut = new StreamFilter(StreamFilterKind.MatchFull);
Assert.Equal(default, sut);
}
[Fact]
public void Should_simplify_input_to_default_filter_with_factory()
{
var sut = StreamFilter.Name();
Assert.Equal(default, sut);
}
}

15
backend/tests/Squidex.Infrastructure.Tests/States/PersistenceBatchTests.cs

@ -185,12 +185,10 @@ public class PersistenceBatchTests
private void SetupEventStore(Dictionary<DomainId, List<MyEvent>> streams) private void SetupEventStore(Dictionary<DomainId, List<MyEvent>> streams)
{ {
var storedStreams = new Dictionary<string, IReadOnlyList<StoredEvent>>(); var events = new List<StoredEvent>();
foreach (var (id, stream) in streams) foreach (var (id, stream) in streams)
{ {
var storedStream = new List<StoredEvent>();
var i = 0; var i = 0;
foreach (var @event in stream) foreach (var @event in stream)
@ -198,23 +196,20 @@ public class PersistenceBatchTests
var eventData = new EventData("Type", new EnvelopeHeaders(), "Payload"); var eventData = new EventData("Type", new EnvelopeHeaders(), "Payload");
var eventStored = new StoredEvent(id.ToString(), i.ToString(CultureInfo.InvariantCulture), i, eventData); var eventStored = new StoredEvent(id.ToString(), i.ToString(CultureInfo.InvariantCulture), i, eventData);
storedStream.Add(eventStored);
A.CallTo(() => eventFormatter.Parse(eventStored)) A.CallTo(() => eventFormatter.Parse(eventStored))
.Returns(new Envelope<IEvent>(@event)); .Returns(new Envelope<IEvent>(@event));
A.CallTo(() => eventFormatter.ParseIfKnown(eventStored)) A.CallTo(() => eventFormatter.ParseIfKnown(eventStored))
.Returns(new Envelope<IEvent>(@event)); .Returns(new Envelope<IEvent>(@event));
events.Add(eventStored);
i++; i++;
} }
storedStreams[id.ToString()] = storedStream;
} }
var streamNames = streams.Keys.Select(x => x.ToString()).ToArray(); var filter = StreamFilter.Name(streams.Keys.Select(x => x.ToString()).ToArray());
A.CallTo(() => eventStore.QueryManyAsync(A<IEnumerable<string>>.That.IsSameSequenceAs(streamNames), A<CancellationToken>._)) A.CallTo(() => eventStore.QueryAllAsync(filter, null, int.MaxValue, A<CancellationToken>._))
.Returns(storedStreams); .Returns(events.ToAsyncEnumerable());
} }
} }

12
backend/tests/Squidex.Infrastructure.Tests/States/PersistenceEventSourcingTests.cs

@ -65,7 +65,7 @@ public class PersistenceEventSourcingTests
{ {
var storedEvent = new StoredEvent("1", "1", 0, new EventData("Type", new EnvelopeHeaders(), "Payload")); var storedEvent = new StoredEvent("1", "1", 0, new EventData("Type", new EnvelopeHeaders(), "Payload"));
A.CallTo(() => eventStore.QueryAsync(key.ToString(), -1, A<CancellationToken>._)) A.CallTo(() => eventStore.QueryStreamAsync(key.ToString(), -1, A<CancellationToken>._))
.Returns(new List<StoredEvent> { storedEvent }); .Returns(new List<StoredEvent> { storedEvent });
A.CallTo(() => eventFormatter.ParseIfKnown(storedEvent)) A.CallTo(() => eventFormatter.ParseIfKnown(storedEvent))
@ -96,7 +96,7 @@ public class PersistenceEventSourcingTests
Assert.False(persistence.IsSnapshotStale); Assert.False(persistence.IsSnapshotStale);
A.CallTo(() => eventStore.QueryAsync(key.ToString(), 2, A<CancellationToken>._)) A.CallTo(() => eventStore.QueryStreamAsync(key.ToString(), 2, A<CancellationToken>._))
.MustHaveHappened(); .MustHaveHappened();
} }
@ -116,7 +116,7 @@ public class PersistenceEventSourcingTests
Assert.True(persistence.IsSnapshotStale); Assert.True(persistence.IsSnapshotStale);
A.CallTo(() => eventStore.QueryAsync(key.ToString(), 1, A<CancellationToken>._)) A.CallTo(() => eventStore.QueryStreamAsync(key.ToString(), 1, A<CancellationToken>._))
.MustHaveHappened(); .MustHaveHappened();
} }
@ -327,7 +327,7 @@ public class PersistenceEventSourcingTests
await persistence.DeleteAsync(); await persistence.DeleteAsync();
A.CallTo(() => eventStore.DeleteStreamAsync(key.ToString(), A<CancellationToken>._)) A.CallTo(() => eventStore.DeleteAsync(StreamFilter.Name(key.ToString()), A<CancellationToken>._))
.MustHaveHappened(); .MustHaveHappened();
A.CallTo(() => snapshotStore.RemoveAsync(key, A<CancellationToken>._)) A.CallTo(() => snapshotStore.RemoveAsync(key, A<CancellationToken>._))
@ -341,7 +341,7 @@ public class PersistenceEventSourcingTests
await persistence.DeleteAsync(); await persistence.DeleteAsync();
A.CallTo(() => eventStore.DeleteStreamAsync(key.ToString(), A<CancellationToken>._)) A.CallTo(() => eventStore.DeleteAsync(StreamFilter.Name(key.ToString()), A<CancellationToken>._))
.MustHaveHappened(); .MustHaveHappened();
A.CallTo(() => snapshotStore.RemoveAsync(key, A<CancellationToken>._)) A.CallTo(() => snapshotStore.RemoveAsync(key, A<CancellationToken>._))
@ -382,7 +382,7 @@ public class PersistenceEventSourcingTests
i++; i++;
} }
A.CallTo(() => eventStore.QueryAsync(key.ToString(), readPosition, A<CancellationToken>._)) A.CallTo(() => eventStore.QueryStreamAsync(key.ToString(), readPosition, A<CancellationToken>._))
.Returns(eventsStored); .Returns(eventsStored);
} }
} }

2
backend/tests/Squidex.Infrastructure.Tests/States/PersistenceSnapshotTests.cs

@ -162,7 +162,7 @@ public class PersistenceSnapshotTests
await persistence.DeleteAsync(); await persistence.DeleteAsync();
A.CallTo(() => eventStore.DeleteStreamAsync(A<string>._, A<CancellationToken>._)) A.CallTo(() => eventStore.DeleteAsync(A<StreamFilter>._, A<CancellationToken>._))
.MustNotHaveHappened(); .MustNotHaveHappened();
A.CallTo(() => snapshotStore.RemoveAsync(key, A<CancellationToken>._)) A.CallTo(() => snapshotStore.RemoveAsync(key, A<CancellationToken>._))

Loading…
Cancel
Save