diff --git a/src/Squidex.Infrastructure.Azure/EventSourcing/Constants.cs b/src/Squidex.Infrastructure.Azure/EventSourcing/Constants.cs new file mode 100644 index 000000000..120f38704 --- /dev/null +++ b/src/Squidex.Infrastructure.Azure/EventSourcing/Constants.cs @@ -0,0 +1,16 @@ +// ========================================================================== +// Squidex Headless CMS +// ========================================================================== +// Copyright (c) Squidex UG (haftungsbeschraenkt) +// All rights reserved. Licensed under the MIT license. +// ========================================================================== + +namespace Squidex.Infrastructure.EventSourcing +{ + internal static class Constants + { + public const string Collection = "Events"; + + public const string LeaseCollection = "Leases"; + } +} diff --git a/src/Squidex.Infrastructure.Azure/EventSourcing/CosmosDbEventStore.cs b/src/Squidex.Infrastructure.Azure/EventSourcing/CosmosDbEventStore.cs index c28984b0f..152bc6bf1 100644 --- a/src/Squidex.Infrastructure.Azure/EventSourcing/CosmosDbEventStore.cs +++ b/src/Squidex.Infrastructure.Azure/EventSourcing/CosmosDbEventStore.cs @@ -11,6 +11,7 @@ using System.Threading; using System.Threading.Tasks; using Microsoft.Azure.Documents; using Microsoft.Azure.Documents.Client; +using Newtonsoft.Json; namespace Squidex.Infrastructure.EventSourcing { @@ -19,27 +20,62 @@ namespace Squidex.Infrastructure.EventSourcing private readonly DocumentClient documentClient; private readonly Uri databaseUri; private readonly Uri collectionUri; + private readonly Uri serviceUri; + private readonly string masterKey; private readonly string databaseId; - private readonly string collectionId; + private readonly JsonSerializerSettings serializerSettings; - public CosmosDbEventStore(DocumentClient documentClient, string database) + public JsonSerializerSettings SerializerSettings { - Guard.NotNull(documentClient, nameof(documentClient)); + get { return serializerSettings; } + } + + public string DatabaseId + { + get { return databaseId; } + } + + public string MasterKey + { + get { return masterKey; } + } + + public Uri ServiceUri + { + get { return serviceUri; } + } + + public CosmosDbEventStore(Uri uri, string masterKey, JsonSerializerSettings serializerSettings, string database) + { + Guard.NotNull(uri, nameof(uri)); + Guard.NotNull(serializerSettings, nameof(serializerSettings)); + Guard.NotNullOrEmpty(masterKey, nameof(masterKey)); Guard.NotNullOrEmpty(database, nameof(database)); - this.documentClient = documentClient; + documentClient = new DocumentClient(uri, masterKey, serializerSettings); databaseUri = UriFactory.CreateDatabaseUri(database); databaseId = database; - collectionUri = UriFactory.CreateDocumentCollectionUri(database, FilterBuilder.Collection); - collectionId = FilterBuilder.Collection; + collectionUri = UriFactory.CreateDocumentCollectionUri(database, Constants.Collection); + + serviceUri = uri; + + this.masterKey = masterKey; + + this.serializerSettings = serializerSettings; } public async Task InitializeAsync(CancellationToken ct = default) { await documentClient.CreateDatabaseIfNotExistsAsync(new Database { Id = databaseId }); + await documentClient.CreateDocumentCollectionIfNotExistsAsync(databaseUri, + new DocumentCollection + { + Id = Constants.LeaseCollection, + }); + await documentClient.CreateDocumentCollectionIfNotExistsAsync(databaseUri, new DocumentCollection { @@ -51,17 +87,17 @@ namespace Squidex.Infrastructure.EventSourcing { Paths = new Collection { - $"/{FilterBuilder.EventStreamField}", - $"/{FilterBuilder.EventStreamOffsetField}" + $"/eventStream", + $"/eventStreamOffset" } } } }, - Id = FilterBuilder.Collection, + Id = Constants.Collection, }, new RequestOptions { - PartitionKey = new PartitionKey($"/{FilterBuilder.EventStreamField}") + PartitionKey = new PartitionKey($"/eventStream") }); } } diff --git a/src/Squidex.Infrastructure.Azure/EventSourcing/CosmosDbEventStore_Reader.cs b/src/Squidex.Infrastructure.Azure/EventSourcing/CosmosDbEventStore_Reader.cs index adaef8785..fb5eb249b 100644 --- a/src/Squidex.Infrastructure.Azure/EventSourcing/CosmosDbEventStore_Reader.cs +++ b/src/Squidex.Infrastructure.Azure/EventSourcing/CosmosDbEventStore_Reader.cs @@ -22,9 +22,8 @@ namespace Squidex.Infrastructure.EventSourcing public IEventSubscription CreateSubscription(IEventSubscriber subscriber, string streamFilter, string position = null) { Guard.NotNull(subscriber, nameof(subscriber)); - Guard.NotNullOrEmpty(streamFilter, nameof(streamFilter)); - throw new NotSupportedException(); + return new CosmosDbSubscription(this, subscriber, streamFilter, position); } public Task CreateIndexAsync(string property) @@ -47,13 +46,13 @@ namespace Squidex.Infrastructure.EventSourcing var commitTimestamp = commit.Timestamp; var commitOffset = 0; - foreach (var e in commit.Events) + foreach (var @event in commit.Events) { eventStreamOffset++; if (eventStreamOffset >= streamPosition) { - var eventData = e.ToEventData(); + var eventData = @event.ToEventData(); var eventToken = new StreamPosition(commitTimestamp, commitOffset, commit.Events.Length); result.Add(new StoredEvent(streamName, eventToken, eventStreamOffset, eventData)); @@ -102,13 +101,13 @@ namespace Squidex.Infrastructure.EventSourcing var commitTimestamp = commit.Timestamp; var commitOffset = 0; - foreach (var e in commit.Events) + foreach (var @event in commit.Events) { eventStreamOffset++; if (commitOffset > lastPosition.CommitOffset || commitTimestamp > lastPosition.Timestamp) { - var eventData = e.ToEventData(); + var eventData = @event.ToEventData(); if (filterExpression(eventData)) { diff --git a/src/Squidex.Infrastructure.Azure/EventSourcing/CosmosDbEventStore_Writer.cs b/src/Squidex.Infrastructure.Azure/EventSourcing/CosmosDbEventStore_Writer.cs index 89650117d..cfcdc2e1e 100644 --- a/src/Squidex.Infrastructure.Azure/EventSourcing/CosmosDbEventStore_Writer.cs +++ b/src/Squidex.Infrastructure.Azure/EventSourcing/CosmosDbEventStore_Writer.cs @@ -21,7 +21,6 @@ namespace Squidex.Infrastructure.EventSourcing { private const int MaxWriteAttempts = 20; private const int MaxCommitSize = 10; - private static readonly FeedOptions TakeOne = new FeedOptions { MaxItemCount = 1 }; public Task DeleteStreamAsync(string streamName) { @@ -29,7 +28,9 @@ namespace Squidex.Infrastructure.EventSourcing return documentClient.QueryAsync(collectionUri, query, commit => { - return documentClient.DeleteDocumentAsync(UriFactory.CreateDocumentUri(databaseId, collectionId, commit.Id.ToString())); + var documentUri = UriFactory.CreateDocumentUri(databaseId, Constants.Collection, commit.Id.ToString()); + + return documentClient.DeleteDocumentAsync(documentUri); }); } diff --git a/src/Squidex.Infrastructure.Azure/EventSourcing/CosmosDbSubscription.cs b/src/Squidex.Infrastructure.Azure/EventSourcing/CosmosDbSubscription.cs new file mode 100644 index 000000000..86fdb41c4 --- /dev/null +++ b/src/Squidex.Infrastructure.Azure/EventSourcing/CosmosDbSubscription.cs @@ -0,0 +1,150 @@ +// ========================================================================== +// Squidex Headless CMS +// ========================================================================== +// Copyright (c) Squidex UG (haftungsbeschraenkt) +// All rights reserved. Licensed under the MIT license. +// ========================================================================== + +using System; +using System.Collections.Generic; +using System.Text.RegularExpressions; +using System.Threading; +using System.Threading.Tasks; +using Microsoft.Azure.Documents; +using Microsoft.Azure.Documents.ChangeFeedProcessor.FeedProcessing; +using Newtonsoft.Json; +using Builder = Microsoft.Azure.Documents.ChangeFeedProcessor.ChangeFeedProcessorBuilder; +using Collection = Microsoft.Azure.Documents.ChangeFeedProcessor.DocumentCollectionInfo; +using Options = Microsoft.Azure.Documents.ChangeFeedProcessor.ChangeFeedProcessorOptions; + +#pragma warning disable IDE0017 // Simplify object initialization + +namespace Squidex.Infrastructure.EventSourcing +{ + internal sealed class CosmosDbSubscription : IEventSubscription, IChangeFeedObserverFactory, IChangeFeedObserver + { + private readonly TaskCompletionSource processorStopRequested = new TaskCompletionSource(); + private readonly Task processorTask; + private readonly CosmosDbEventStore store; + private readonly Regex regex; + private readonly string hostName; + private readonly IEventSubscriber subscriber; + + public CosmosDbSubscription(CosmosDbEventStore store, IEventSubscriber subscriber, string streamFilter, string position = null) + { + var fromBeginning = string.IsNullOrWhiteSpace(position); + + if (fromBeginning) + { + hostName = $"squidex.{DateTime.UtcNow.Ticks.ToString()}"; + } + else + { + hostName = position; + } + + if (!StreamFilter.IsAll(streamFilter)) + { + regex = new Regex(streamFilter); + } + + this.store = store; + + this.subscriber = subscriber; + + processorTask = Task.Run(async () => + { + try + { + Collection CreateCollection(string name) + { + var collection = new Collection(); + + collection.CollectionName = name; + collection.DatabaseName = store.DatabaseId; + collection.MasterKey = store.MasterKey; + collection.Uri = store.ServiceUri; + + return collection; + } + + var processor = + await new Builder() + .WithFeedCollection(CreateCollection(Constants.Collection)) + .WithLeaseCollection(CreateCollection(Constants.LeaseCollection)) + .WithHostName(hostName) + .WithProcessorOptions(new Options { StartFromBeginning = fromBeginning }) + .WithObserverFactory(this) + .BuildAsync(); + + await processor.StartAsync(); + await processorStopRequested.Task; + await processor.StopAsync(); + } + catch (Exception ex) + { + await subscriber.OnErrorAsync(this, ex); + } + }); + } + + public IChangeFeedObserver CreateObserver() + { + return this; + } + + public async Task CloseAsync(IChangeFeedObserverContext context, ChangeFeedObserverCloseReason reason) + { + if (reason == ChangeFeedObserverCloseReason.ObserverError) + { + await subscriber.OnErrorAsync(this, new InvalidOperationException("Change feed observer failed.")); + } + } + + public Task OpenAsync(IChangeFeedObserverContext context) + { + return Task.CompletedTask; + } + + public async Task ProcessChangesAsync(IChangeFeedObserverContext context, IReadOnlyList docs, CancellationToken cancellationToken) + { + if (!processorStopRequested.Task.IsCompleted) + { + foreach (var document in docs) + { + if (!processorStopRequested.Task.IsCompleted) + { + var streamName = document.GetPropertyValue("eventStream"); + + if (regex == null || regex.IsMatch(streamName)) + { + var commit = JsonConvert.DeserializeObject(document.ToString(), store.SerializerSettings); + + var eventStreamOffset = (int)commit.EventStreamOffset; + + foreach (var @event in commit.Events) + { + eventStreamOffset++; + + var eventData = @event.ToEventData(); + + await subscriber.OnEventAsync(this, new StoredEvent(commit.EventStream, hostName, eventStreamOffset, eventData)); + } + } + } + } + } + } + + public void WakeUp() + { + } + + public Task StopAsync() + { + processorStopRequested.SetResult(true); + + return processorTask; + } + } +} diff --git a/src/Squidex.Infrastructure.Azure/EventSourcing/FilterBuilder.cs b/src/Squidex.Infrastructure.Azure/EventSourcing/FilterBuilder.cs index 821c53092..7180ad59c 100644 --- a/src/Squidex.Infrastructure.Azure/EventSourcing/FilterBuilder.cs +++ b/src/Squidex.Infrastructure.Azure/EventSourcing/FilterBuilder.cs @@ -18,14 +18,6 @@ namespace Squidex.Infrastructure.EventSourcing { internal static class FilterBuilder { - public const string Collection = "Events"; - - public static readonly string CommitId = nameof(CosmosDbEventCommit.Id).ToCamelCase(); - public static readonly string EventsCountField = nameof(CosmosDbEventCommit.EventsCount).ToCamelCase(); - public static readonly string EventStreamOffsetField = nameof(CosmosDbEventCommit.EventStreamOffset).ToCamelCase(); - public static readonly string EventStreamField = nameof(CosmosDbEventCommit.EventStream).ToCamelCase(); - public static readonly string TimestampField = nameof(CosmosDbEventCommit.Timestamp).ToCamelCase(); - public static async Task QueryAsync(this DocumentClient documentClient, Uri collectionUri, SqlQuerySpec querySpec, Func handler, CancellationToken ct = default) { var query = @@ -52,12 +44,12 @@ namespace Squidex.Infrastructure.EventSourcing { var query = $"SELECT TOP 1 " + - $" e.{CommitId}," + - $" e.{EventsCountField} " + - $"FROM {Collection} e " + + $" e.id," + + $" e.eventsCount " + + $"FROM {Constants.Collection} e " + $"WHERE " + - $" e.{EventStreamField} = @name " + - $"ORDER BY e.{EventStreamOffsetField} DESC"; + $" e.eventStream = @name " + + $"ORDER BY e.eventStreamOffset DESC"; var parameters = new SqlParameterCollection { @@ -71,12 +63,12 @@ namespace Squidex.Infrastructure.EventSourcing { var query = $"SELECT TOP 1 " + - $" e.{EventStreamOffsetField}," + - $" e.{EventsCountField} " + - $"FROM {Collection} e " + + $" e.eventStreamOffset," + + $" e.eventsCount " + + $"FROM {Constants.Collection} e " + $"WHERE " + - $" e.{EventStreamField} = @name " + - $"ORDER BY e.{EventStreamOffsetField} DESC"; + $" e.eventStream = @name " + + $"ORDER BY e.eventStreamOffset DESC"; var parameters = new SqlParameterCollection { @@ -90,11 +82,11 @@ namespace Squidex.Infrastructure.EventSourcing { var query = $"SELECT * " + - $"FROM {Collection} e " + + $"FROM {Constants.Collection} e " + $"WHERE " + - $" e.{EventStreamField} = @name " + - $"AND e.{EventStreamOffsetField} >= @position " + - $"ORDER BY e.{EventStreamOffsetField} ASC"; + $" e.eventStream = @name " + + $"AND e.eventStreamOffset >= @position " + + $"ORDER BY e.eventStreamOffset ASC"; var parameters = new SqlParameterCollection { @@ -131,7 +123,7 @@ namespace Squidex.Infrastructure.EventSourcing private static SqlQuerySpec BuildQuery(List filters, SqlParameterCollection parameters) { - var query = $"SELECT * FROM {Collection} e WHERE {string.Join(" AND ", filters)} ORDER BY e.{TimestampField}"; + var query = $"SELECT * FROM {Constants.Collection} e WHERE {string.Join(" AND ", filters)} ORDER BY e.timestamp"; return new SqlQuerySpec(query, parameters); } @@ -145,15 +137,15 @@ namespace Squidex.Infrastructure.EventSourcing private static void ForRegex(this List filters, SqlParameterCollection parameters, string streamFilter) { - if (!string.IsNullOrWhiteSpace(streamFilter) && !string.Equals(streamFilter, ".*", StringComparison.OrdinalIgnoreCase)) + if (!StreamFilter.IsAll(streamFilter)) { if (streamFilter.Contains("^")) { - filters.Add($"STARTSWITH(e.{EventStreamField}, @filter)"); + filters.Add($"STARTSWITH(e.eventStream, @filter)"); } else { - filters.Add($"e.{EventStreamField} = @filter"); + filters.Add($"e.eventStream = @filter"); } parameters.Add(new SqlParameter("@filter", streamFilter)); @@ -164,11 +156,11 @@ namespace Squidex.Infrastructure.EventSourcing { if (streamPosition.IsEndOfCommit) { - filters.Add($"e.{TimestampField} > @time"); + filters.Add($"e.timestamp > @time"); } else { - filters.Add($"e.{TimestampField} >= @time"); + filters.Add($"e.timestamp >= @time"); } parameters.Add(new SqlParameter("@time", streamPosition.Timestamp)); diff --git a/src/Squidex.Infrastructure.Azure/Squidex.Infrastructure.Azure.csproj b/src/Squidex.Infrastructure.Azure/Squidex.Infrastructure.Azure.csproj index d4d6148da..32ea4d0ad 100644 --- a/src/Squidex.Infrastructure.Azure/Squidex.Infrastructure.Azure.csproj +++ b/src/Squidex.Infrastructure.Azure/Squidex.Infrastructure.Azure.csproj @@ -5,6 +5,7 @@ 7.3 + diff --git a/src/Squidex.Infrastructure.MongoDb/EventSourcing/MongoEventStore_Reader.cs b/src/Squidex.Infrastructure.MongoDb/EventSourcing/MongoEventStore_Reader.cs index 314304445..e0e7d635c 100644 --- a/src/Squidex.Infrastructure.MongoDb/EventSourcing/MongoEventStore_Reader.cs +++ b/src/Squidex.Infrastructure.MongoDb/EventSourcing/MongoEventStore_Reader.cs @@ -55,13 +55,13 @@ namespace Squidex.Infrastructure.EventSourcing var commitTimestamp = commit.Timestamp; var commitOffset = 0; - foreach (var e in commit.Events) + foreach (var @event in commit.Events) { eventStreamOffset++; if (eventStreamOffset >= streamPosition) { - var eventData = e.ToEventData(); + var eventData = @event.ToEventData(); var eventToken = new StreamPosition(commitTimestamp, commitOffset, commit.Events.Length); result.Add(new StoredEvent(streamName, eventToken, eventStreamOffset, eventData)); @@ -108,13 +108,13 @@ namespace Squidex.Infrastructure.EventSourcing var commitTimestamp = commit.Timestamp; var commitOffset = 0; - foreach (var e in commit.Events) + foreach (var @event in commit.Events) { eventStreamOffset++; if (commitOffset > lastPosition.CommitOffset || commitTimestamp > lastPosition.Timestamp) { - var eventData = e.ToEventData(); + var eventData = @event.ToEventData(); if (filterExpression(eventData)) { @@ -157,7 +157,7 @@ namespace Squidex.Infrastructure.EventSourcing private static void AppendByStream(string streamFilter, List filters) { - if (!string.IsNullOrWhiteSpace(streamFilter) && !string.Equals(streamFilter, ".*", StringComparison.OrdinalIgnoreCase)) + if (!StreamFilter.IsAll(streamFilter)) { if (streamFilter.Contains("^")) { diff --git a/src/Squidex.Infrastructure/EventSourcing/StreamFilter.cs b/src/Squidex.Infrastructure/EventSourcing/StreamFilter.cs new file mode 100644 index 000000000..b3bc063af --- /dev/null +++ b/src/Squidex.Infrastructure/EventSourcing/StreamFilter.cs @@ -0,0 +1,22 @@ +// ========================================================================== +// Squidex Headless CMS +// ========================================================================== +// Copyright (c) Squidex UG (haftungsbeschraenkt) +// All rights reserved. Licensed under the MIT license. +// ========================================================================== + +using System; + +namespace Squidex.Infrastructure.EventSourcing +{ + public static class StreamFilter + { + public static bool IsAll(string filter) + { + return string.IsNullOrWhiteSpace(filter) + || string.Equals(filter, ".*", StringComparison.OrdinalIgnoreCase) + || string.Equals(filter, "(.*)", StringComparison.OrdinalIgnoreCase) + || string.Equals(filter, "(.*?)", StringComparison.OrdinalIgnoreCase); + } + } +} diff --git a/tests/Squidex.Infrastructure.Tests/EventSourcing/CosmosDbEventStoreFixture.cs b/tests/Squidex.Infrastructure.Tests/EventSourcing/CosmosDbEventStoreFixture.cs index 221a6c305..42a5292be 100644 --- a/tests/Squidex.Infrastructure.Tests/EventSourcing/CosmosDbEventStoreFixture.cs +++ b/tests/Squidex.Infrastructure.Tests/EventSourcing/CosmosDbEventStoreFixture.cs @@ -13,20 +13,23 @@ namespace Squidex.Infrastructure.EventSourcing { public sealed class CosmosDbEventStoreFixture : IDisposable { + private const string EmulatorKey = "C2y6yDjf5/R+ob0N8A7Cgv30VRDJIWEHLM+4QDU5DE2nQ9nDuVTqobD4b8mGGyPMbIZnqyMsEcaGQy67XIw/Jw=="; + private const string EmulatorUri = "https://localhost:8081"; private readonly DocumentClient client; public CosmosDbEventStore EventStore { get; } public CosmosDbEventStoreFixture() { - client = new DocumentClient(new Uri("https://localhost:8081"), "C2y6yDjf5/R+ob0N8A7Cgv30VRDJIWEHLM+4QDU5DE2nQ9nDuVTqobD4b8mGGyPMbIZnqyMsEcaGQy67XIw/Jw==", JsonHelper.DefaultSettings()); + client = new DocumentClient(new Uri(EmulatorUri), EmulatorKey); - EventStore = new CosmosDbEventStore(client, "Test"); + EventStore = new CosmosDbEventStore(new Uri(EmulatorUri), EmulatorKey, JsonHelper.DefaultSettings(), "Test"); EventStore.InitializeAsync().Wait(); } public void Dispose() { + client.DeleteDatabaseAsync(UriFactory.CreateDatabaseUri("Test")).Wait(); } } }