diff --git a/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Assets/MongoAssetFolderRepository_SnapshotStore.cs b/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Assets/MongoAssetFolderRepository_SnapshotStore.cs index 9319377f8..23eca83ed 100644 --- a/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Assets/MongoAssetFolderRepository_SnapshotStore.cs +++ b/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Assets/MongoAssetFolderRepository_SnapshotStore.cs @@ -55,9 +55,9 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Assets { using (Telemetry.Activities.StartActivity("MongoAssetFolderRepository/WriteAsync")) { - var entity = MongoAssetFolderEntity.Create(job); + var entityJob = job.As(MongoAssetFolderEntity.Create(job)); - await Collection.UpsertVersionedAsync(job.Key, job.OldVersion, job.NewVersion, entity, ct); + await Collection.UpsertVersionedAsync(entityJob, ct); } } diff --git a/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Assets/MongoAssetRepository_SnapshotStore.cs b/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Assets/MongoAssetRepository_SnapshotStore.cs index 0127ed81b..5623a8ebe 100644 --- a/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Assets/MongoAssetRepository_SnapshotStore.cs +++ b/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Assets/MongoAssetRepository_SnapshotStore.cs @@ -55,9 +55,9 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Assets { using (Telemetry.Activities.StartActivity("MongoAssetRepository/WriteAsync")) { - var entity = MongoAssetEntity.Create(job); + var entityJob = job.As(MongoAssetEntity.Create(job)); - await Collection.UpsertVersionedAsync(job.Key, job.OldVersion, job.NewVersion, entity, ct); + await Collection.UpsertVersionedAsync(entityJob, ct); } } diff --git a/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Contents/MongoContentCollection.cs b/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Contents/MongoContentCollection.cs index 4b6d2b551..51fe183c6 100644 --- a/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Contents/MongoContentCollection.cs +++ b/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Contents/MongoContentCollection.cs @@ -16,6 +16,7 @@ using Squidex.Domain.Apps.Entities.Schemas; using Squidex.Infrastructure; using Squidex.Infrastructure.MongoDb; using Squidex.Infrastructure.Queries; +using Squidex.Infrastructure.States; using Squidex.Infrastructure.Translations; #pragma warning disable IDE0060 // Remove unused parameter @@ -251,25 +252,47 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Contents return Collection.Find(FindAll).ToAsyncEnumerable(ct); } - public async Task UpsertVersionedAsync(DomainId documentId, long oldVersion, MongoContentEntity value, + public async Task UpsertAsync(SnapshotWriteJob job, CancellationToken ct = default) { if (queryInDedicatedCollection != null) { - await queryInDedicatedCollection.UpsertVersionedAsync(documentId, oldVersion, value, default); + await queryInDedicatedCollection.UpsertAsync(job, ct); } - await Collection.UpsertVersionedAsync(documentId, oldVersion, value.Version, value, default); + await Collection.ReplaceOneAsync(Filter.Eq(x => x.DocumentId, job.Key), job.Value, UpsertReplace, ct); + } + + public async Task UpsertVersionedAsync(IClientSessionHandle session, SnapshotWriteJob job, + CancellationToken ct = default) + { + if (queryInDedicatedCollection != null) + { + await queryInDedicatedCollection.UpsertVersionedAsync(session, job, ct); + } + + await Collection.UpsertVersionedAsync(session, job, ct); } public async Task RemoveAsync(DomainId key, CancellationToken ct = default) { - var previous = await Collection.FindOneAndDeleteAsync(x => x.DocumentId == key, null, default); + var previous = await Collection.FindOneAndDeleteAsync(x => x.DocumentId == key, null, ct); + + if (queryInDedicatedCollection != null && previous != null) + { + await queryInDedicatedCollection.RemoveAsync(previous, ct); + } + } + + public async Task RemoveAsync(IClientSessionHandle session, DomainId key, + CancellationToken ct = default) + { + var previous = await Collection.FindOneAndDeleteAsync(session, x => x.DocumentId == key, null, ct); if (queryInDedicatedCollection != null && previous != null) { - await queryInDedicatedCollection.RemoveAsync(previous, default); + await queryInDedicatedCollection.RemoveAsync(session, previous, ct); } } diff --git a/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Contents/MongoContentRepository.cs b/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Contents/MongoContentRepository.cs index 810b8f6a3..65c9ffd1b 100644 --- a/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Contents/MongoContentRepository.cs +++ b/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Contents/MongoContentRepository.cs @@ -7,6 +7,7 @@ using Microsoft.Extensions.Options; using MongoDB.Driver; +using MongoDB.Driver.Core.Clusters; using NodaTime; using Squidex.Domain.Apps.Core.Contents; using Squidex.Domain.Apps.Entities.Apps; @@ -24,8 +25,11 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Contents { private readonly MongoContentCollection collectionComplete; private readonly MongoContentCollection collectionPublished; + private readonly IMongoDatabase database; private readonly IAppProvider appProvider; + public bool CanUseTransactions { get; private set; } + static MongoContentRepository() { BsonStringSerializer.Register(); @@ -34,6 +38,8 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Contents public MongoContentRepository(IMongoDatabase database, IAppProvider appProvider, IOptions options) { + this.database = database; + collectionComplete = new MongoContentCollection("States_Contents_All3", database, ReadPreference.Primary, options.Value.OptimizeForSelfHosting); @@ -50,6 +56,11 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Contents { await collectionComplete.InitializeAsync(ct); await collectionPublished.InitializeAsync(ct); + + var clusterVersion = await database.GetMajorVersionAsync(ct); + var clusteredAsReplica = database.Client.Cluster.Description.Type == ClusterType.ReplicaSet; + + CanUseTransactions = clusteredAsReplica && clusterVersion >= 4; } public IAsyncEnumerable StreamAll(DomainId appId, HashSet? schemaIds, diff --git a/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Contents/MongoContentRepository_SnapshotStore.cs b/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Contents/MongoContentRepository_SnapshotStore.cs index c6ba30681..a5f5f65f3 100644 --- a/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Contents/MongoContentRepository_SnapshotStore.cs +++ b/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Contents/MongoContentRepository_SnapshotStore.cs @@ -33,6 +33,7 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Contents var existing = await collectionComplete.FindAsync(key, ct); + // Support for all versions, where we do not have full snapshots in the collection. if (existing?.IsSnapshot == true) { return new SnapshotResult(existing.DocumentId, existing.ToState(), existing.Version); @@ -67,6 +68,7 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Contents { using (Telemetry.Activities.StartActivity("MongoContentRepository/RemoveAsync")) { + // Some data is corrupt and might throw an exception if we do not ignore it. if (key == DomainId.Empty) { return; @@ -83,14 +85,33 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Contents { using (Telemetry.Activities.StartActivity("MongoContentRepository/WriteAsync")) { + // Some data is corrupt and might throw an exception if we do not ignore it. if (!IsValid(job.Value)) { return; } - await Task.WhenAll( - UpsertFrontendAsync(job, ct), - UpsertPublishedAsync(job, ct)); + if (!CanUseTransactions) + { + // If transactions are not supported we update the documents without version checks, + // otherwise we would not be able to recover from inconsistencies. + await Task.WhenAll( + UpsertCompleteAsync(job, default), + UpsertPublishedAsync(job, default)); + return; + } + + using (var session = await database.Client.StartSessionAsync(cancellationToken: ct)) + { + // Make an update with full transaction support to be more consistent. + await session.WithTransactionAsync(async (session, ct) => + { + await Task.WhenAll( + UpsertVersionedCompleteAsync(session, job, ct), + UpsertVersionedPublishedAsync(session, job, ct)); + return true; + }, null, ct); + } } } @@ -99,11 +120,11 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Contents { using (Telemetry.Activities.StartActivity("MongoContentRepository/WriteManyAsync")) { - var updates = new Dictionary, List>(); + var collectionUpdates = new Dictionary, List>(); var add = new Action, MongoContentEntity>((collection, entity) => { - updates.GetOrAddNew(collection).Add(entity); + collectionUpdates.GetOrAddNew(collection).Add(entity); }); foreach (var job in jobs) @@ -123,7 +144,15 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Contents } } - await Parallel.ForEachAsync(updates, ct, (update, ct) => + var parallelOptions = new ParallelOptions + { + CancellationToken = ct, + // This is just an estimate, but we do not want ot have unlimited parallelism. + MaxDegreeOfParallelism = 8 + }; + + // Make one update per collection. + await Parallel.ForEachAsync(collectionUpdates, parallelOptions, (update, ct) => { return new ValueTask(update.Key.InsertManyAsync(update.Value, InsertUnordered, ct)); }); @@ -131,38 +160,49 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Contents } private async Task UpsertPublishedAsync(SnapshotWriteJob job, - CancellationToken ct = default) + CancellationToken ct) { if (ShouldWritePublished(job.Value)) { - await UpsertPublishedContentAsync(job, ct); + var entityJob = job.As(await MongoContentEntity.CreatePublishedAsync(job, appProvider)); + + await collectionPublished.UpsertAsync(entityJob, ct); } else { - await DeletePublishedContentAsync(job.Value.UniqueId, ct); + await collectionPublished.RemoveAsync(job.Key, ct); } } - private Task DeletePublishedContentAsync(DomainId key, - CancellationToken ct = default) + private async Task UpsertVersionedPublishedAsync(IClientSessionHandle session, SnapshotWriteJob job, + CancellationToken ct) { - return collectionPublished.RemoveAsync(key, ct); + if (ShouldWritePublished(job.Value)) + { + var entityJob = job.As(await MongoContentEntity.CreatePublishedAsync(job, appProvider)); + + await collectionPublished.UpsertVersionedAsync(session, entityJob, ct); + } + else + { + await collectionPublished.RemoveAsync(session, job.Key, ct); + } } - private async Task UpsertFrontendAsync(SnapshotWriteJob job, - CancellationToken ct = default) + private async Task UpsertCompleteAsync(SnapshotWriteJob job, + CancellationToken ct) { - var entity = await MongoContentEntity.CreateAsync(job, appProvider); + var entityJob = job.As(await MongoContentEntity.CreateAsync(job, appProvider)); - await collectionComplete.UpsertVersionedAsync(entity.DocumentId, job.OldVersion, entity, ct); + await collectionComplete.UpsertAsync(entityJob, ct); } - private async Task UpsertPublishedContentAsync(SnapshotWriteJob job, - CancellationToken ct = default) + private async Task UpsertVersionedCompleteAsync(IClientSessionHandle session, SnapshotWriteJob job, + CancellationToken ct) { - var entity = await MongoContentEntity.CreatePublishedAsync(job, appProvider); + var entityJob = job.As(await MongoContentEntity.CreateAsync(job, appProvider)); - await collectionPublished.UpsertVersionedAsync(entity.DocumentId, job.OldVersion, entity, ct); + await collectionComplete.UpsertVersionedAsync(session, entityJob, ct); } private static bool ShouldWritePublished(ContentDomainObject.State value) diff --git a/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Contents/Operations/QueryInDedicatedCollection.cs b/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Contents/Operations/QueryInDedicatedCollection.cs index 998b110b2..ca664a1ea 100644 --- a/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Contents/Operations/QueryInDedicatedCollection.cs +++ b/backend/src/Squidex.Domain.Apps.Entities.MongoDb/Contents/Operations/QueryInDedicatedCollection.cs @@ -15,6 +15,7 @@ using Squidex.Infrastructure; using Squidex.Infrastructure.MongoDb; using Squidex.Infrastructure.MongoDb.Queries; using Squidex.Infrastructure.Queries; +using Squidex.Infrastructure.States; namespace Squidex.Domain.Apps.Entities.MongoDb.Contents.Operations { @@ -110,12 +111,20 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Contents.Operations return ResultList.Create(contentTotal, contentEntities); } - public async Task UpsertVersionedAsync(DomainId documentId, long oldVersion, MongoContentEntity value, + public async Task UpsertAsync(SnapshotWriteJob job, CancellationToken ct = default) { - var collection = await GetCollectionAsync(value.AppId.Id, value.SchemaId.Id); + var collection = await GetCollectionAsync(job.Value.AppId.Id, job.Value.SchemaId.Id); + + await collection.ReplaceOneAsync(Filter.Eq(x => x.DocumentId, job.Key), job.Value, UpsertReplace, ct); + } + + public async Task UpsertVersionedAsync(IClientSessionHandle session, SnapshotWriteJob job, + CancellationToken ct = default) + { + var collection = await GetCollectionAsync(job.Value.AppId.Id, job.Value.SchemaId.Id); - await collection.UpsertVersionedAsync(documentId, oldVersion, value.Version, value, ct); + await collection.UpsertVersionedAsync(session, job, ct); } public async Task RemoveAsync(MongoContentEntity value, @@ -126,6 +135,14 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Contents.Operations await collection.DeleteOneAsync(x => x.DocumentId == value.DocumentId, null, ct); } + public async Task RemoveAsync(IClientSessionHandle session, MongoContentEntity value, + CancellationToken ct = default) + { + var collection = await GetCollectionAsync(value.AppId.Id, value.SchemaId.Id); + + await collection.DeleteOneAsync(session, x => x.DocumentId == value.DocumentId, null, ct); + } + private static FilterDefinition BuildFilter(FilterNode? filter) { var filters = new List> diff --git a/backend/src/Squidex.Domain.Apps.Entities/Assets/DomainObject/AssetCommandMiddleware.cs b/backend/src/Squidex.Domain.Apps.Entities/Assets/DomainObject/AssetCommandMiddleware.cs index 5a6e9b1f0..2f6634515 100644 --- a/backend/src/Squidex.Domain.Apps.Entities/Assets/DomainObject/AssetCommandMiddleware.cs +++ b/backend/src/Squidex.Domain.Apps.Entities/Assets/DomainObject/AssetCommandMiddleware.cs @@ -99,8 +99,6 @@ namespace Squidex.Domain.Apps.Entities.Assets.DomainObject finally { await assetFileStore.DeleteAsync(tempFile, ct); - - await command.File.DisposeAsync(); } } @@ -119,8 +117,6 @@ namespace Squidex.Domain.Apps.Entities.Assets.DomainObject finally { await assetFileStore.DeleteAsync(tempFile, ct); - - await command.File.DisposeAsync(); } } diff --git a/backend/src/Squidex.Infrastructure.MongoDb/MongoDb/MongoExtensions.cs b/backend/src/Squidex.Infrastructure.MongoDb/MongoDb/MongoExtensions.cs index eef6bf902..0326fdb5c 100644 --- a/backend/src/Squidex.Infrastructure.MongoDb/MongoDb/MongoExtensions.cs +++ b/backend/src/Squidex.Infrastructure.MongoDb/MongoDb/MongoExtensions.cs @@ -119,21 +119,60 @@ namespace Squidex.Infrastructure.MongoDb } } - public static async Task UpsertVersionedAsync(this IMongoCollection collection, TKey key, long oldVersion, long newVersion, T document, + public static async Task UpsertVersionedAsync(this IMongoCollection collection, IClientSessionHandle session, SnapshotWriteJob job, CancellationToken ct = default) - where T : IVersionedEntity where TKey : notnull + where T : IVersionedEntity { + var (key, snapshot, newVersion, oldVersion) = job; try { - document.DocumentId = key; - document.Version = newVersion; + snapshot.DocumentId = key; + snapshot.Version = newVersion; Expression> filter = oldVersion > EtagVersion.Any ? x => x.DocumentId.Equals(key) && x.Version == oldVersion : x => x.DocumentId.Equals(key); - var result = await collection.ReplaceOneAsync(filter, document, UpsertReplace, ct); + var result = await collection.ReplaceOneAsync(session, filter, job.Value, UpsertReplace, ct); + + return result.IsAcknowledged && result.ModifiedCount == 1; + } + catch (MongoWriteException ex) when (ex.WriteError?.Category == ServerErrorCategory.DuplicateKey) + { + var existingVersion = + await collection.Find(session, x => x.DocumentId.Equals(key)).Only(x => x.DocumentId, x => x.Version) + .FirstOrDefaultAsync(ct); + + if (existingVersion != null) + { + var field = Field.Of(x => nameof(x.Version)); + + throw new InconsistentStateException(existingVersion[field].AsInt64, oldVersion); + } + else + { + throw new InconsistentStateException(EtagVersion.Any, oldVersion); + } + } + } + + public static async Task UpsertVersionedAsync(this IMongoCollection collection, SnapshotWriteJob job, + CancellationToken ct = default) + where T : IVersionedEntity + { + var (key, snapshot, newVersion, oldVersion) = job; + try + { + snapshot.DocumentId = key; + snapshot.Version = newVersion; + + Expression> filter = + oldVersion > EtagVersion.Any ? + x => x.DocumentId.Equals(key) && x.Version == oldVersion : + x => x.DocumentId.Equals(key); + + var result = await collection.ReplaceOneAsync(filter, snapshot, UpsertReplace, ct); return result.IsAcknowledged && result.ModifiedCount == 1; } diff --git a/backend/src/Squidex.Infrastructure.MongoDb/States/MongoSnapshotStoreBase.cs b/backend/src/Squidex.Infrastructure.MongoDb/States/MongoSnapshotStoreBase.cs index 37c0f3da4..7c1b7d3d5 100644 --- a/backend/src/Squidex.Infrastructure.MongoDb/States/MongoSnapshotStoreBase.cs +++ b/backend/src/Squidex.Infrastructure.MongoDb/States/MongoSnapshotStoreBase.cs @@ -51,9 +51,9 @@ namespace Squidex.Infrastructure.States { using (Telemetry.Activities.StartActivity("MongoSnapshotStoreBase/WriteAsync")) { - var document = CreateDocument(job.Key, job.Value, job.OldVersion); + var entityJob = job.As(CreateDocument(job.Key, job.Value, job.OldVersion)); - await Collection.UpsertVersionedAsync(job.Key, job.OldVersion, job.NewVersion, document, ct); + await Collection.UpsertVersionedAsync(entityJob, ct); } } diff --git a/backend/src/Squidex.Infrastructure/States/ISnapshotStore.cs b/backend/src/Squidex.Infrastructure/States/ISnapshotStore.cs index 4bb0d24ad..6192dba87 100644 --- a/backend/src/Squidex.Infrastructure/States/ISnapshotStore.cs +++ b/backend/src/Squidex.Infrastructure/States/ISnapshotStore.cs @@ -33,5 +33,11 @@ namespace Squidex.Infrastructure.States public record struct SnapshotResult(DomainId Key, T Value, long Version, bool IsValid = true); - public record struct SnapshotWriteJob(DomainId Key, T Value, long NewVersion, long OldVersion = EtagVersion.Any); + public record struct SnapshotWriteJob(DomainId Key, T Value, long NewVersion, long OldVersion = EtagVersion.Any) + { + public SnapshotWriteJob As(TOther snapshot) + { + return new SnapshotWriteJob(Key, snapshot, NewVersion, OldVersion); + } + } } diff --git a/backend/src/Squidex/Config/Domain/LoggingServices.cs b/backend/src/Squidex/Config/Domain/LoggingServices.cs index 72ddd7e6a..16e0e261c 100644 --- a/backend/src/Squidex/Config/Domain/LoggingServices.cs +++ b/backend/src/Squidex/Config/Domain/LoggingServices.cs @@ -18,11 +18,11 @@ namespace Squidex.Config.Domain public static void ConfigureForSquidex(this ILoggingBuilder builder, IConfiguration config) { builder.ClearProviders(); + + // Also adds semantic logging. builder.ConfigureSemanticLog(config); builder.AddConfiguration(config.GetSection("logging")); - - builder.AddSemanticLog(); builder.AddFilters(); builder.Services.AddServices(config); diff --git a/backend/tools/TestSuite/TestSuite.ApiTests/ContentUpdateTests.cs b/backend/tools/TestSuite/TestSuite.ApiTests/ContentUpdateTests.cs index 501155105..137f4463b 100644 --- a/backend/tools/TestSuite/TestSuite.ApiTests/ContentUpdateTests.cs +++ b/backend/tools/TestSuite/TestSuite.ApiTests/ContentUpdateTests.cs @@ -373,49 +373,52 @@ namespace TestSuite.ApiTests [Fact] public async Task Should_update_content_in_parallel() { - TestEntity content = null; - try + for (var i = 0; i < 200; i++) { - // STEP 1: Create a new item. - content = await _.Contents.CreateAsync(new TestEntityData { Number = 2 }, ContentCreateOptions.AsPublish); + TestEntity content = null; + try + { + // STEP 1: Create a new item. + content = await _.Contents.CreateAsync(new TestEntityData { Number = 2 }, ContentCreateOptions.AsPublish); - var numErrors = 0; - var numSuccess = 0; + var numErrors = 0; + var numSuccess = 0; - // STEP 3: Make parallel updates. - await Parallel.ForEachAsync(Enumerable.Range(0, 20), async (i, ct) => - { - try + // STEP 3: Make parallel updates. + await Parallel.ForEachAsync(Enumerable.Range(0, 20), async (i, ct) => { - await _.Contents.UpdateAsync(content.Id, new TestEntityData { Number = i }); + try + { + await _.Contents.UpdateAsync(content.Id, new TestEntityData { Number = i }); - Interlocked.Increment(ref numSuccess); - } - catch (SquidexException ex) when (ex.StatusCode is 409 or 412) - { - Interlocked.Increment(ref numErrors); - return; - } - }); + Interlocked.Increment(ref numSuccess); + } + catch (SquidexException ex) when (ex.StatusCode is 409 or 412) + { + Interlocked.Increment(ref numErrors); + return; + } + }); - // At least some errors and success should have happened. - Assert.True(numErrors > 0); - Assert.True(numSuccess > 0); + // At least some errors and success should have happened. + Assert.True(numErrors > 0); + Assert.True(numSuccess > 0); - // STEP 3: Make an normal update to ensure nothing is corrupt. - await _.Contents.UpdateAsync(content.Id, new TestEntityData { Number = 2 }); + // STEP 3: Make an normal update to ensure nothing is corrupt. + await _.Contents.UpdateAsync(content.Id, new TestEntityData { Number = 2 }); - var updated = await _.Contents.GetAsync(content.Id); + var updated = await _.Contents.GetAsync(content.Id); - Assert.Equal(2, content.Data.Number); - } - finally - { - if (content != null) + Assert.Equal(2, content.Data.Number); + } + finally { - await _.Contents.DeleteAsync(content.Id); + if (content != null) + { + await _.Contents.DeleteAsync(content.Id); + } } } }