From 191b29c1c737501c9d6f6674c93170c5fbbe6d4a Mon Sep 17 00:00:00 2001 From: Sebastian Date: Wed, 25 Jul 2018 10:58:48 +0200 Subject: [PATCH] Extracted restore behavior. --- .../Backup/BackupGrain.cs | 25 ++- .../Backup/EventStreamReader.cs | 10 +- .../Backup/EventStreamWriter.cs | 12 +- .../Backup/Handlers/HandlerBase.cs | 55 +++++ .../Backup/Handlers/RestoreApp.cs | 50 +++++ .../Backup/Handlers/RestoreAssets.cs | 72 +++++++ .../Backup/Handlers/RestoreContents.cs | 54 +++++ .../Backup/Handlers/RestoreRules.cs | 70 ++++++ .../Backup/Handlers/RestoreSchemas.cs | 75 +++++++ .../Backup/IRestoreGrain.cs | 21 ++ .../Backup/IRestoreHandler.cs | 24 +++ .../Backup/IRestoreJob.cs | 23 ++ .../Backup/RestoreGrain.cs | 203 +++++++++++++++++- .../Backup/State/RestoreState.cs | 17 ++ .../Backup/State/RestoreStateJob.cs | 38 ++++ .../Tags/GrainTagService.cs | 5 + .../Tags/ITagGrain.cs | 1 + .../Tags/ITagService.cs | 2 + .../Tags/TagGrain.cs | 5 + .../EventSourcing/Formatter.cs | 1 + .../EventSourcing/MongoEventStore_Reader.cs | 4 +- .../EventSourcing/StoredEvent.cs | 7 +- .../Backup/EventStreamTests.cs | 12 +- .../Grains/EventConsumerGrainTests.cs | 12 +- .../EventSourcing/RetrySubscriptionTests.cs | 4 +- .../States/PersistenceEventSourcingTests.cs | 4 +- 26 files changed, 759 insertions(+), 47 deletions(-) create mode 100644 src/Squidex.Domain.Apps.Entities/Backup/Handlers/HandlerBase.cs create mode 100644 src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreApp.cs create mode 100644 src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreAssets.cs create mode 100644 src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreContents.cs create mode 100644 src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreRules.cs create mode 100644 src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreSchemas.cs create mode 100644 src/Squidex.Domain.Apps.Entities/Backup/IRestoreGrain.cs create mode 100644 src/Squidex.Domain.Apps.Entities/Backup/IRestoreHandler.cs create mode 100644 src/Squidex.Domain.Apps.Entities/Backup/IRestoreJob.cs create mode 100644 src/Squidex.Domain.Apps.Entities/Backup/State/RestoreState.cs create mode 100644 src/Squidex.Domain.Apps.Entities/Backup/State/RestoreStateJob.cs diff --git a/src/Squidex.Domain.Apps.Entities/Backup/BackupGrain.cs b/src/Squidex.Domain.Apps.Entities/Backup/BackupGrain.cs index 5b2346bab..12c601806 100644 --- a/src/Squidex.Domain.Apps.Entities/Backup/BackupGrain.cs +++ b/src/Squidex.Domain.Apps.Entities/Backup/BackupGrain.cs @@ -11,7 +11,6 @@ using System.Linq; using System.Threading; using System.Threading.Tasks; using NodaTime; -using Orleans; using Orleans.Concurrency; using Squidex.Domain.Apps.Entities.Backup.State; using Squidex.Domain.Apps.Events; @@ -191,28 +190,28 @@ namespace Squidex.Domain.Apps.Entities.Backup var assetVersion = 0L; var assetId = Guid.Empty; - if (parsedEvent.Payload is AssetCreated assetCreated) + switch (parsedEvent.Payload) { - assetId = assetCreated.AssetId; - assetVersion = assetCreated.FileVersion; + case AssetCreated assetCreated: + assetId = assetCreated.AssetId; + assetVersion = assetCreated.FileVersion; + break; + case AssetUpdated asetUpdated: + assetId = asetUpdated.AssetId; + assetVersion = asetUpdated.FileVersion; + break; } - if (parsedEvent.Payload is AssetUpdated asetUpdated) + await writer.WriteEventAsync(@event, attachment => { - assetId = asetUpdated.AssetId; - assetVersion = asetUpdated.FileVersion; - } - - await writer.WriteEventAsync(eventData, async attachmentStream => - { - await assetStore.DownloadAsync(assetId.ToString(), assetVersion, null, attachmentStream); + return assetStore.DownloadAsync(assetId.ToString(), assetVersion, null, attachment); }); job.HandledAssets++; } else { - await writer.WriteEventAsync(eventData); + await writer.WriteEventAsync(@event); } job.HandledEvents++; diff --git a/src/Squidex.Domain.Apps.Entities/Backup/EventStreamReader.cs b/src/Squidex.Domain.Apps.Entities/Backup/EventStreamReader.cs index ba1d4bcbf..836accdff 100644 --- a/src/Squidex.Domain.Apps.Entities/Backup/EventStreamReader.cs +++ b/src/Squidex.Domain.Apps.Entities/Backup/EventStreamReader.cs @@ -34,7 +34,7 @@ namespace Squidex.Domain.Apps.Entities.Backup } } - public async Task ReadEventsAsync(Func eventHandler) + public async Task ReadEventsAsync(Func eventHandler) { Guard.NotNull(eventHandler, nameof(eventHandler)); @@ -52,18 +52,16 @@ namespace Squidex.Domain.Apps.Entities.Backup break; } - EventData eventData; + StoredEvent eventData; using (var stream = eventEntry.Open()) { using (var textReader = new StreamReader(stream)) { - eventData = (EventData)JsonSerializer.Deserialize(textReader, typeof(EventData)); + eventData = (StoredEvent)JsonSerializer.Deserialize(textReader, typeof(StoredEvent)); } } - readEvents++; - var attachmentFolder = readAttachments / MaxItemsPerFolder; var attachmentPath = $"attachments/{attachmentFolder}/{readEvents}.blob"; var attachmentEntry = archive.GetEntry(attachmentPath); @@ -81,6 +79,8 @@ namespace Squidex.Domain.Apps.Entities.Backup { await eventHandler(eventData, null); } + + readEvents++; } } } diff --git a/src/Squidex.Domain.Apps.Entities/Backup/EventStreamWriter.cs b/src/Squidex.Domain.Apps.Entities/Backup/EventStreamWriter.cs index 27ea2f9f2..2ca38b729 100644 --- a/src/Squidex.Domain.Apps.Entities/Backup/EventStreamWriter.cs +++ b/src/Squidex.Domain.Apps.Entities/Backup/EventStreamWriter.cs @@ -36,13 +36,9 @@ namespace Squidex.Domain.Apps.Entities.Backup } } - public async Task WriteEventAsync(EventData eventData, Func attachment = null) + public async Task WriteEventAsync(StoredEvent storedEvent, Func attachment = null) { - var eventObject = - new JObject( - new JProperty("type", eventData.Type), - new JProperty("payload", eventData.Payload), - new JProperty("metadata", eventData.Metadata)); + var eventObject = JObject.FromObject(storedEvent); var eventFolder = writtenEvents / MaxItemsPerFolder; var eventPath = $"events/{eventFolder}/{writtenEvents}.json"; @@ -59,8 +55,6 @@ namespace Squidex.Domain.Apps.Entities.Backup } } - writtenEvents++; - if (attachment != null) { var attachmentFolder = writtenAttachments / MaxItemsPerFolder; @@ -74,6 +68,8 @@ namespace Squidex.Domain.Apps.Entities.Backup writtenAttachments++; } + + writtenEvents++; } } } diff --git a/src/Squidex.Domain.Apps.Entities/Backup/Handlers/HandlerBase.cs b/src/Squidex.Domain.Apps.Entities/Backup/Handlers/HandlerBase.cs new file mode 100644 index 000000000..c43a91a10 --- /dev/null +++ b/src/Squidex.Domain.Apps.Entities/Backup/Handlers/HandlerBase.cs @@ -0,0 +1,55 @@ +// ========================================================================== +// Squidex Headless CMS +// ========================================================================== +// Copyright (c) Squidex UG (haftungsbeschraenkt) +// All rights reserved. Licensed under the MIT license. +// ========================================================================== + +using System; +using System.Collections.Generic; +using System.Threading.Tasks; +using Squidex.Infrastructure; +using Squidex.Infrastructure.Commands; +using Squidex.Infrastructure.EventSourcing; +using Squidex.Infrastructure.States; + +namespace Squidex.Domain.Apps.Entities.Backup.Handlers +{ + public abstract class HandlerBase + { + private readonly IStore store; + + protected HandlerBase(IStore store) + { + Guard.NotNull(store, nameof(store)); + + this.store = store; + } + + protected async Task RebuildManyAsync(IEnumerable ids, Func action) + { + foreach (var id in ids) + { + await action(id); + } + } + + protected async Task RebuildAsync(Guid key, Func, TState, TState> func) where TState : IDomainState, new() + { + var state = new TState + { + Version = EtagVersion.Empty + }; + + var persistence = store.WithSnapshotsAndEventSourcing(typeof(TGrain), key, s => state = s, e => + { + state = func(e, state); + + state.Version++; + }); + + await persistence.ReadAsync(); + await persistence.WriteSnapshotAsync(state); + } + } +} diff --git a/src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreApp.cs b/src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreApp.cs new file mode 100644 index 000000000..77a65aabb --- /dev/null +++ b/src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreApp.cs @@ -0,0 +1,50 @@ +// ========================================================================== +// Squidex Headless CMS +// ========================================================================== +// Copyright (c) Squidex UG (haftungsbeschraenkt) +// All rights reserved. Licensed under the MIT license. +// ========================================================================== + +using System; +using System.IO; +using System.Threading.Tasks; +using Squidex.Domain.Apps.Events.Apps; +using Squidex.Infrastructure; +using Squidex.Infrastructure.EventSourcing; +using Squidex.Infrastructure.States; +using Squidex.Infrastructure.Tasks; + +namespace Squidex.Domain.Apps.Entities.Backup.Handlers +{ + public sealed class RestoreApp : HandlerBase, IRestoreHandler + { + private NamedId appId; + + public string Name { get; } = "App"; + + public RestoreApp(IStore store) + : base(store) + { + } + + public Task HandleAsync(Envelope @event, Stream attachment) + { + if (@event.Payload is AppCreated appCreated) + { + appId = appCreated.AppId; + } + + return TaskHelper.Done; + } + + public Task ProcessAsync() + { + return TaskHelper.Done; + } + + public Task CompleteAsync() + { + return TaskHelper.Done; + } + } +} diff --git a/src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreAssets.cs b/src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreAssets.cs new file mode 100644 index 000000000..6b16947bc --- /dev/null +++ b/src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreAssets.cs @@ -0,0 +1,72 @@ +// ========================================================================== +// Squidex Headless CMS +// ========================================================================== +// Copyright (c) Squidex UG (haftungsbeschraenkt) +// All rights reserved. Licensed under the MIT license. +// ========================================================================== + +using System; +using System.Collections.Generic; +using System.IO; +using System.Threading.Tasks; +using Squidex.Domain.Apps.Entities.Assets; +using Squidex.Domain.Apps.Entities.Assets.State; +using Squidex.Domain.Apps.Events.Assets; +using Squidex.Infrastructure; +using Squidex.Infrastructure.Assets; +using Squidex.Infrastructure.EventSourcing; +using Squidex.Infrastructure.States; +using Squidex.Infrastructure.Tasks; + +namespace Squidex.Domain.Apps.Entities.Backup.Handlers +{ + public sealed class RestoreAssets : HandlerBase, IRestoreHandler + { + private readonly HashSet assetIds = new HashSet(); + private readonly IAssetStore assetStore; + + public string Name { get; } = "Assets"; + + public RestoreAssets(IStore store, IAssetStore assetStore) + : base(store) + { + Guard.NotNull(assetStore, nameof(assetStore)); + + this.assetStore = assetStore; + } + + public async Task HandleAsync(Envelope @event, Stream attachment) + { + var assetVersion = 0L; + var assetId = Guid.Empty; + + switch (@event.Payload) + { + case AssetCreated assetCreated: + assetId = assetCreated.AssetId; + assetVersion = assetCreated.FileVersion; + assetIds.Add(assetCreated.AssetId); + break; + case AssetUpdated asetUpdated: + assetId = asetUpdated.AssetId; + assetVersion = asetUpdated.FileVersion; + break; + } + + if (attachment != null) + { + await assetStore.UploadAsync(assetId.ToString(), assetVersion, null, attachment); + } + } + + public Task ProcessAsync() + { + return RebuildManyAsync(assetIds, id => RebuildAsync(id, (e, s) => s.Apply(e))); + } + + public Task CompleteAsync() + { + return TaskHelper.Done; + } + } +} diff --git a/src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreContents.cs b/src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreContents.cs new file mode 100644 index 000000000..e5d7be6fd --- /dev/null +++ b/src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreContents.cs @@ -0,0 +1,54 @@ +// ========================================================================== +// Squidex Headless CMS +// ========================================================================== +// Copyright (c) Squidex UG (haftungsbeschraenkt) +// All rights reserved. Licensed under the MIT license. +// ========================================================================== + +using System; +using System.Collections.Generic; +using System.IO; +using System.Threading.Tasks; +using Squidex.Domain.Apps.Entities.Contents; +using Squidex.Domain.Apps.Entities.Contents.State; +using Squidex.Domain.Apps.Events.Contents; +using Squidex.Infrastructure.EventSourcing; +using Squidex.Infrastructure.States; +using Squidex.Infrastructure.Tasks; + +namespace Squidex.Domain.Apps.Entities.Backup.Handlers +{ + public sealed class RestoreContents : HandlerBase, IRestoreHandler + { + private readonly HashSet contentIds = new HashSet(); + + public string Name { get; } = "Contents"; + + public RestoreContents(IStore store) + : base(store) + { + } + + public Task HandleAsync(Envelope @event, Stream attachment) + { + switch (@event.Payload) + { + case ContentCreated contentCreated: + contentIds.Add(contentCreated.ContentId); + break; + } + + return TaskHelper.Done; + } + + public Task ProcessAsync() + { + return RebuildManyAsync(contentIds, id => RebuildAsync(id, (e, s) => s.Apply(e))); + } + + public Task CompleteAsync() + { + return TaskHelper.Done; + } + } +} diff --git a/src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreRules.cs b/src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreRules.cs new file mode 100644 index 000000000..44a280f78 --- /dev/null +++ b/src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreRules.cs @@ -0,0 +1,70 @@ +// ========================================================================== +// Squidex Headless CMS +// ========================================================================== +// Copyright (c) Squidex UG (haftungsbeschraenkt) +// All rights reserved. Licensed under the MIT license. +// ========================================================================== + +using System; +using System.Collections.Generic; +using System.IO; +using System.Threading.Tasks; +using Orleans; +using Squidex.Domain.Apps.Entities.Rules; +using Squidex.Domain.Apps.Entities.Rules.State; +using Squidex.Domain.Apps.Events.Apps; +using Squidex.Domain.Apps.Events.Rules; +using Squidex.Infrastructure; +using Squidex.Infrastructure.EventSourcing; +using Squidex.Infrastructure.States; +using Squidex.Infrastructure.Tasks; + +namespace Squidex.Domain.Apps.Entities.Backup.Handlers +{ + public sealed class RestoreRules : HandlerBase, IRestoreHandler + { + private readonly HashSet ruleIds = new HashSet(); + private readonly IGrainFactory grainFactory; + private Guid appId; + + public string Name { get; } = "Rules"; + + public RestoreRules(IStore store, IGrainFactory grainFactory) + : base(store) + { + Guard.NotNull(grainFactory, nameof(grainFactory)); + + this.grainFactory = grainFactory; + } + + public Task HandleAsync(Envelope @event, Stream attachment) + { + switch (@event.Payload) + { + case AppCreated appCreated: + appId = appCreated.AppId.Id; + break; + case RuleCreated ruleCreated: + ruleIds.Add(ruleCreated.RuleId); + break; + case RuleDeleted ruleDeleted: + ruleIds.Remove(ruleDeleted.RuleId); + break; + } + + return TaskHelper.Done; + } + + public async Task ProcessAsync() + { + await RebuildManyAsync(ruleIds, id => RebuildAsync(id, (e, s) => s.Apply(e))); + + await grainFactory.GetGrain(appId).RebuildAsync(ruleIds); + } + + public Task CompleteAsync() + { + return TaskHelper.Done; + } + } +} diff --git a/src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreSchemas.cs b/src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreSchemas.cs new file mode 100644 index 000000000..c56924e8b --- /dev/null +++ b/src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreSchemas.cs @@ -0,0 +1,75 @@ +// ========================================================================== +// Squidex Headless CMS +// ========================================================================== +// Copyright (c) Squidex UG (haftungsbeschraenkt) +// All rights reserved. Licensed under the MIT license. +// ========================================================================== + +using System; +using System.Collections.Generic; +using System.IO; +using System.Linq; +using System.Threading.Tasks; +using Orleans; +using Squidex.Domain.Apps.Core.Schemas; +using Squidex.Domain.Apps.Entities.Schemas; +using Squidex.Domain.Apps.Entities.Schemas.State; +using Squidex.Domain.Apps.Events.Apps; +using Squidex.Domain.Apps.Events.Schemas; +using Squidex.Infrastructure; +using Squidex.Infrastructure.EventSourcing; +using Squidex.Infrastructure.States; +using Squidex.Infrastructure.Tasks; + +namespace Squidex.Domain.Apps.Entities.Backup.Handlers +{ + public sealed class RestoreSchemas : HandlerBase, IRestoreHandler + { + private readonly HashSet> schemaIds = new HashSet>(); + private readonly Dictionary schemasByName = new Dictionary(); + private readonly FieldRegistry fieldRegistry; + private readonly IGrainFactory grainFactory; + private Guid appId; + + public string Name { get; } = "Schemas"; + + public RestoreSchemas(IStore store, FieldRegistry fieldRegistry, IGrainFactory grainFactory) + : base(store) + { + Guard.NotNull(fieldRegistry, nameof(fieldRegistry)); + Guard.NotNull(grainFactory, nameof(grainFactory)); + + this.fieldRegistry = fieldRegistry; + + this.grainFactory = grainFactory; + } + + public Task HandleAsync(Envelope @event, Stream attachment) + { + switch (@event.Payload) + { + case AppCreated appCreated: + appId = appCreated.AppId.Id; + break; + case SchemaCreated schemaCreated: + schemaIds.Add(schemaCreated.SchemaId); + schemasByName[schemaCreated.SchemaId.Name] = schemaCreated.SchemaId.Id; + break; + } + + return TaskHelper.Done; + } + + public async Task ProcessAsync() + { + await RebuildManyAsync(schemaIds.Select(x => x.Id), id => RebuildAsync(id, (e, s) => s.Apply(e, fieldRegistry))); + + await grainFactory.GetGrain(appId).RebuildAsync(schemasByName); + } + + public Task CompleteAsync() + { + return TaskHelper.Done; + } + } +} diff --git a/src/Squidex.Domain.Apps.Entities/Backup/IRestoreGrain.cs b/src/Squidex.Domain.Apps.Entities/Backup/IRestoreGrain.cs new file mode 100644 index 000000000..219169dd3 --- /dev/null +++ b/src/Squidex.Domain.Apps.Entities/Backup/IRestoreGrain.cs @@ -0,0 +1,21 @@ +// ========================================================================== +// Squidex Headless CMS +// ========================================================================== +// Copyright (c) Squidex UG (haftungsbeschraenkt) +// All rights reserved. Licensed under the MIT license. +// ========================================================================== + +using System; +using System.Threading.Tasks; +using Squidex.Infrastructure; +using Squidex.Infrastructure.Orleans; + +namespace Squidex.Domain.Apps.Entities.Backup +{ + public interface IRestoreGrain + { + Task RestoreAsync(Uri url, RefToken user); + + Task> GetStateAsync(); + } +} diff --git a/src/Squidex.Domain.Apps.Entities/Backup/IRestoreHandler.cs b/src/Squidex.Domain.Apps.Entities/Backup/IRestoreHandler.cs new file mode 100644 index 000000000..ed943c250 --- /dev/null +++ b/src/Squidex.Domain.Apps.Entities/Backup/IRestoreHandler.cs @@ -0,0 +1,24 @@ +// ========================================================================== +// Squidex Headless CMS +// ========================================================================== +// Copyright (c) Squidex UG (haftungsbeschraenkt) +// All rights reserved. Licensed under the MIT license. +// ========================================================================== + +using System.IO; +using System.Threading.Tasks; +using Squidex.Infrastructure.EventSourcing; + +namespace Squidex.Domain.Apps.Entities.Backup +{ + public interface IRestoreHandler + { + string Name { get; } + + Task HandleAsync(Envelope @event, Stream attachment); + + Task ProcessAsync(); + + Task CompleteAsync(); + } +} diff --git a/src/Squidex.Domain.Apps.Entities/Backup/IRestoreJob.cs b/src/Squidex.Domain.Apps.Entities/Backup/IRestoreJob.cs new file mode 100644 index 000000000..b3714d15d --- /dev/null +++ b/src/Squidex.Domain.Apps.Entities/Backup/IRestoreJob.cs @@ -0,0 +1,23 @@ +// ========================================================================== +// Squidex Headless CMS +// ========================================================================== +// Copyright (c) Squidex UG (haftungsbeschraenkt) +// All rights reserved. Licensed under the MIT license. +// ========================================================================== + +using System; +using NodaTime; + +namespace Squidex.Domain.Apps.Entities.Backup +{ + public interface IRestoreJob + { + Uri Uri { get; } + + Instant Started { get; } + + bool IsFailed { get; } + + string Status { get; } + } +} diff --git a/src/Squidex.Domain.Apps.Entities/Backup/RestoreGrain.cs b/src/Squidex.Domain.Apps.Entities/Backup/RestoreGrain.cs index eb1d5757b..524405df2 100644 --- a/src/Squidex.Domain.Apps.Entities/Backup/RestoreGrain.cs +++ b/src/Squidex.Domain.Apps.Entities/Backup/RestoreGrain.cs @@ -5,14 +5,213 @@ // All rights reserved. Licensed under the MIT license. // ========================================================================== +using System; +using System.Collections.Generic; +using System.Net.Http; +using System.Threading.Tasks; +using NodaTime; +using Orleans; +using Squidex.Domain.Apps.Entities.Backup.State; +using Squidex.Domain.Apps.Events; +using Squidex.Infrastructure; +using Squidex.Infrastructure.Assets; +using Squidex.Infrastructure.EventSourcing; +using Squidex.Infrastructure.Log; using Squidex.Infrastructure.Orleans; +using Squidex.Infrastructure.States; +using Squidex.Infrastructure.Tasks; namespace Squidex.Domain.Apps.Entities.Backup { - public sealed class RestoreGrain : GrainOfString + public sealed class RestoreGrain : GrainOfString, IRestoreGrain { - public RestoreGrain() + private static readonly Duration UpdateDuration = Duration.FromSeconds(1); + private readonly IClock clock; + private readonly IAssetStore assetStore; + private readonly IEventDataFormatter eventDataFormatter; + private readonly IAppCleanerGrain appCleaner; + private readonly ISemanticLog log; + private readonly IEventStore eventStore; + private readonly IBackupArchiveLocation backupArchiveLocation; + private readonly IStore store; + private readonly IEnumerable handlers; + private RestoreState state = new RestoreState(); + private IPersistence persistence; + + public RestoreGrain( + IAssetStore assetStore, + IBackupArchiveLocation backupArchiveLocation, + IClock clock, + IEventStore eventStore, + IEventDataFormatter eventDataFormatter, + IGrainFactory grainFactory, + IEnumerable handlers, + ISemanticLog log, + IStore store) + { + Guard.NotNull(assetStore, nameof(assetStore)); + Guard.NotNull(backupArchiveLocation, nameof(backupArchiveLocation)); + Guard.NotNull(clock, nameof(clock)); + Guard.NotNull(eventStore, nameof(eventStore)); + Guard.NotNull(eventDataFormatter, nameof(eventDataFormatter)); + Guard.NotNull(grainFactory, nameof(grainFactory)); + Guard.NotNull(handlers, nameof(handlers)); + Guard.NotNull(store, nameof(store)); + Guard.NotNull(log, nameof(log)); + + this.assetStore = assetStore; + this.backupArchiveLocation = backupArchiveLocation; + this.clock = clock; + this.eventStore = eventStore; + this.eventDataFormatter = eventDataFormatter; + this.handlers = handlers; + this.store = store; + this.log = log; + } + + public override async Task OnActivateAsync(string key) + { + persistence = store.WithSnapshots(GetType(), key, s => state = s); + + await persistence.ReadAsync(); + + await CleanupAsync(); + } + + public Task RestoreAsync(Uri url, RefToken user) + { + if (state.Job != null) + { + throw new DomainException("A restore operation is already running."); + } + + state.Job = new RestoreStateJob { Started = clock.GetCurrentInstant(), Uri = url, User = user }; + + return ProcessAsync(); + } + + private async Task CleanupAsync() + { + if (state.Job != null) + { + state.Job.Status = "Failed due application restart"; + state.Job.IsFailed = true; + + if (state.Job.AppId != Guid.Empty) + { + appCleaner.EnqueueAppAsync(state.Job.AppId).Forget(); + } + + await persistence.WriteSnapshotAsync(state); + } + } + + private async Task ProcessAsync() + { + try + { + await DoAsync( + "Downloading Backup", + "Downloaded Backup", + DownloadAsync); + + await DoAsync( + "Reading Events", + "Readed Events", + ReadEventsAsync); + + foreach (var handler in handlers) + { + await DoAsync($"{handler.Name} Proessing", $"{handler.Name} Processed", handler.ProcessAsync); + } + + foreach (var handler in handlers) + { + await DoAsync($"{handler.Name} Completing", $"{handler.Name} Completed", handler.CompleteAsync); + } + + state.Job = null; + } + catch (Exception ex) + { + log.LogError(ex, w => w + .WriteProperty("action", "makeBackup") + .WriteProperty("status", "failed") + .WriteProperty("backupId", state.Job.Id.ToString())); + + state.Job.IsFailed = true; + + if (state.Job.AppId != Guid.Empty) + { + appCleaner.EnqueueAppAsync(state.Job.AppId).Forget(); + } + } + finally + { + await persistence.WriteSnapshotAsync(state); + } + } + + private async Task DownloadAsync() + { + using (var client = new HttpClient()) + { + using (var sourceStream = await client.GetStreamAsync(state.Job.Uri.ToString())) + { + using (var targetStream = await backupArchiveLocation.OpenStreamAsync(state.Job.Id)) + { + await sourceStream.CopyToAsync(targetStream); + } + } + } + } + + private async Task ReadEventsAsync() + { + using (var stream = await backupArchiveLocation.OpenStreamAsync(state.Job.Id)) + { + using (var reader = new EventStreamReader(stream)) + { + var eventIndex = 0; + + await reader.ReadEventsAsync(async (@event, attachment) => + { + var eventData = @event.Data; + + var parsedEvent = eventDataFormatter.Parse(eventData); + + if (parsedEvent.Payload is SquidexEvent squidexEvent) + { + squidexEvent.Actor = state.Job.User; + } + + foreach (var handler in handlers) + { + await handler.HandleAsync(parsedEvent, attachment); + } + + await eventStore.AppendAsync(Guid.NewGuid(), @event.StreamName, new List { @event.Data }); + + eventIndex++; + + state.Job.Status = $"Handled event {eventIndex}"; + }); + } + } + } + + private async Task DoAsync(string start, string end, Func action) + { + state.Job.Status = start; + + await action(); + + state.Job.Status = end; + } + + public Task> GetStateAsync() { + return Task.FromResult>(state.Job); } } } diff --git a/src/Squidex.Domain.Apps.Entities/Backup/State/RestoreState.cs b/src/Squidex.Domain.Apps.Entities/Backup/State/RestoreState.cs new file mode 100644 index 000000000..d86a658fe --- /dev/null +++ b/src/Squidex.Domain.Apps.Entities/Backup/State/RestoreState.cs @@ -0,0 +1,17 @@ +// ========================================================================== +// Squidex Headless CMS +// ========================================================================== +// Copyright (c) Squidex UG (haftungsbeschraenkt) +// All rights reserved. Licensed under the MIT license. +// ========================================================================== + +using Newtonsoft.Json; + +namespace Squidex.Domain.Apps.Entities.Backup.State +{ + public class RestoreState + { + [JsonProperty] + public RestoreStateJob Job { get; set; } + } +} diff --git a/src/Squidex.Domain.Apps.Entities/Backup/State/RestoreStateJob.cs b/src/Squidex.Domain.Apps.Entities/Backup/State/RestoreStateJob.cs new file mode 100644 index 000000000..9c794bae3 --- /dev/null +++ b/src/Squidex.Domain.Apps.Entities/Backup/State/RestoreStateJob.cs @@ -0,0 +1,38 @@ +// ========================================================================== +// Squidex Headless CMS +// ========================================================================== +// Copyright (c) Squidex UG (haftungsbeschraenkt) +// All rights reserved. Licensed under the MIT license. +// ========================================================================== + +using System; +using Newtonsoft.Json; +using NodaTime; +using Squidex.Infrastructure; + +namespace Squidex.Domain.Apps.Entities.Backup.State +{ + public sealed class RestoreStateJob : IRestoreJob + { + [JsonProperty] + public Guid Id { get; set; } = Guid.NewGuid(); + + [JsonProperty] + public Guid AppId { get; set; } + + [JsonProperty] + public RefToken User { get; set; } + + [JsonProperty] + public Uri Uri { get; set; } + + [JsonProperty] + public Instant Started { get; set; } + + [JsonProperty] + public string Status { get; set; } + + [JsonProperty] + public bool IsFailed { get; set; } + } +} diff --git a/src/Squidex.Domain.Apps.Entities/Tags/GrainTagService.cs b/src/Squidex.Domain.Apps.Entities/Tags/GrainTagService.cs index 22f6cb017..490189f50 100644 --- a/src/Squidex.Domain.Apps.Entities/Tags/GrainTagService.cs +++ b/src/Squidex.Domain.Apps.Entities/Tags/GrainTagService.cs @@ -49,6 +49,11 @@ namespace Squidex.Domain.Apps.Entities.Tags return GetGrain(appId, group).GetTagsAsync(); } + public Task RebuildTagsAsync(Guid appId, string group, Dictionary allTags) + { + return GetGrain(appId, group).RebuildTagsAsync(allTags); + } + public Task ClearAsync(Guid appId) { return GetGrain(appId, TagGroups.Assets).ClearAsync(); diff --git a/src/Squidex.Domain.Apps.Entities/Tags/ITagGrain.cs b/src/Squidex.Domain.Apps.Entities/Tags/ITagGrain.cs index 2947aebda..8218255c8 100644 --- a/src/Squidex.Domain.Apps.Entities/Tags/ITagGrain.cs +++ b/src/Squidex.Domain.Apps.Entities/Tags/ITagGrain.cs @@ -22,5 +22,6 @@ namespace Squidex.Domain.Apps.Entities.Tags Task> GetTagsAsync(); Task ClearAsync(); + Task RebuildTagsAsync(Dictionary allTags); } } diff --git a/src/Squidex.Domain.Apps.Entities/Tags/ITagService.cs b/src/Squidex.Domain.Apps.Entities/Tags/ITagService.cs index b4972b99c..370921b42 100644 --- a/src/Squidex.Domain.Apps.Entities/Tags/ITagService.cs +++ b/src/Squidex.Domain.Apps.Entities/Tags/ITagService.cs @@ -20,5 +20,7 @@ namespace Squidex.Domain.Apps.Entities.Tags Task> DenormalizeTagsAsync(Guid appId, string group, HashSet ids); Task> GetTagsAsync(Guid appId, string group); + + Task RebuildTagsAsync(Guid appId, string group, Dictionary allTags); } } diff --git a/src/Squidex.Domain.Apps.Entities/Tags/TagGrain.cs b/src/Squidex.Domain.Apps.Entities/Tags/TagGrain.cs index c0f9c6fac..a40110ac4 100644 --- a/src/Squidex.Domain.Apps.Entities/Tags/TagGrain.cs +++ b/src/Squidex.Domain.Apps.Entities/Tags/TagGrain.cs @@ -58,6 +58,11 @@ namespace Squidex.Domain.Apps.Entities.Tags return persistence.DeleteAsync(); } + public Task RebuildTagsAsync(Dictionary allTags) + { + throw new NotImplementedException(); + } + public async Task> NormalizeTagsAsync(HashSet names, HashSet ids) { var result = new HashSet(); diff --git a/src/Squidex.Infrastructure.GetEventStore/EventSourcing/Formatter.cs b/src/Squidex.Infrastructure.GetEventStore/EventSourcing/Formatter.cs index aefd883b2..c28abf358 100644 --- a/src/Squidex.Infrastructure.GetEventStore/EventSourcing/Formatter.cs +++ b/src/Squidex.Infrastructure.GetEventStore/EventSourcing/Formatter.cs @@ -24,6 +24,7 @@ namespace Squidex.Infrastructure.EventSourcing var eventData = new EventData { Type = @event.EventType, Payload = body, Metadata = meta }; return new StoredEvent( + @event.EventStreamId, resolvedEvent.OriginalEventNumber.ToString(), resolvedEvent.Event.EventNumber, eventData); diff --git a/src/Squidex.Infrastructure.MongoDb/EventSourcing/MongoEventStore_Reader.cs b/src/Squidex.Infrastructure.MongoDb/EventSourcing/MongoEventStore_Reader.cs index 07861115d..a9e31eb43 100644 --- a/src/Squidex.Infrastructure.MongoDb/EventSourcing/MongoEventStore_Reader.cs +++ b/src/Squidex.Infrastructure.MongoDb/EventSourcing/MongoEventStore_Reader.cs @@ -60,7 +60,7 @@ namespace Squidex.Infrastructure.EventSourcing var eventData = e.ToEventData(); var eventToken = new StreamPosition(commitTimestamp, commitOffset, commit.Events.Length); - result.Add(new StoredEvent(eventToken, eventStreamOffset, eventData)); + result.Add(new StoredEvent(streamName, eventToken, eventStreamOffset, eventData)); } } } @@ -111,7 +111,7 @@ namespace Squidex.Infrastructure.EventSourcing var eventData = e.ToEventData(); var eventToken = new StreamPosition(commitTimestamp, commitOffset, commit.Events.Length); - await callback(new StoredEvent(eventToken, eventStreamOffset, eventData)); + await callback(new StoredEvent(commit.EventStream, eventToken, eventStreamOffset, eventData)); commitOffset++; } diff --git a/src/Squidex.Infrastructure/EventSourcing/StoredEvent.cs b/src/Squidex.Infrastructure/EventSourcing/StoredEvent.cs index 3c93e21a4..97ee0f55c 100644 --- a/src/Squidex.Infrastructure/EventSourcing/StoredEvent.cs +++ b/src/Squidex.Infrastructure/EventSourcing/StoredEvent.cs @@ -9,14 +9,17 @@ namespace Squidex.Infrastructure.EventSourcing { public sealed class StoredEvent { + public string StreamName { get; } + public string EventPosition { get; } public long EventStreamNumber { get; } public EventData Data { get; } - public StoredEvent(string eventPosition, long eventStreamNumber, EventData data) + public StoredEvent(string streamName, string eventPosition, long eventStreamNumber, EventData data) { + Guard.NotNullOrEmpty(streamName, nameof(streamName)); Guard.NotNullOrEmpty(eventPosition, nameof(eventPosition)); Guard.NotNull(data, nameof(data)); @@ -24,6 +27,8 @@ namespace Squidex.Infrastructure.EventSourcing EventPosition = eventPosition; EventStreamNumber = eventStreamNumber; + + StreamName = streamName; } } } diff --git a/tests/Squidex.Domain.Apps.Entities.Tests/Backup/EventStreamTests.cs b/tests/Squidex.Domain.Apps.Entities.Tests/Backup/EventStreamTests.cs index 816daa499..6cadd083f 100644 --- a/tests/Squidex.Domain.Apps.Entities.Tests/Backup/EventStreamTests.cs +++ b/tests/Squidex.Domain.Apps.Entities.Tests/Backup/EventStreamTests.cs @@ -18,7 +18,7 @@ namespace Squidex.Domain.Apps.Entities.Backup { public sealed class EventInfo { - public EventData Data { get; set; } + public StoredEvent Stored { get; set; } public byte[] Attachment { get; set; } } @@ -33,7 +33,7 @@ namespace Squidex.Domain.Apps.Entities.Backup for (var i = 0; i < 1000; i++) { var eventData = new EventData { Type = i.ToString(), Metadata = i, Payload = i }; - var eventInfo = new EventInfo { Data = eventData }; + var eventInfo = new EventInfo { Stored = new StoredEvent("S", "1", 2, eventData) }; if (i % 10 == 0) { @@ -49,11 +49,11 @@ namespace Squidex.Domain.Apps.Entities.Backup { if (@event.Attachment == null) { - await reader.WriteEventAsync(@event.Data); + await reader.WriteEventAsync(@event.Stored); } else { - await reader.WriteEventAsync(@event.Data, s => s.WriteAsync(@event.Attachment, 0, 1)); + await reader.WriteEventAsync(@event.Stored, s => s.WriteAsync(@event.Attachment, 0, 1)); } } } @@ -64,9 +64,9 @@ namespace Squidex.Domain.Apps.Entities.Backup using (var reader = new EventStreamReader(stream)) { - await reader.ReadEventsAsync(async (eventData, attachment) => + await reader.ReadEventsAsync(async (stored, attachment) => { - var eventInfo = new EventInfo { Data = eventData }; + var eventInfo = new EventInfo { Stored = stored }; if (attachment != null) { diff --git a/tests/Squidex.Infrastructure.Tests/EventSourcing/Grains/EventConsumerGrainTests.cs b/tests/Squidex.Infrastructure.Tests/EventSourcing/Grains/EventConsumerGrainTests.cs index d6556c694..da62e3e23 100644 --- a/tests/Squidex.Infrastructure.Tests/EventSourcing/Grains/EventConsumerGrainTests.cs +++ b/tests/Squidex.Infrastructure.Tests/EventSourcing/Grains/EventConsumerGrainTests.cs @@ -173,7 +173,7 @@ namespace Squidex.Infrastructure.EventSourcing.Grains [Fact] public async Task Should_invoke_and_update_position_when_event_received() { - var @event = new StoredEvent(Guid.NewGuid().ToString(), 123, eventData); + var @event = new StoredEvent("Stream", Guid.NewGuid().ToString(), 123, eventData); await sut.OnActivateAsync(consumerName); await sut.ActivateAsync(); @@ -195,7 +195,7 @@ namespace Squidex.Infrastructure.EventSourcing.Grains A.CallTo(() => formatter.Parse(eventData, true)) .Throws(new TypeNameNotFoundException()); - var @event = new StoredEvent(Guid.NewGuid().ToString(), 123, eventData); + var @event = new StoredEvent("Stream", Guid.NewGuid().ToString(), 123, eventData); await sut.OnActivateAsync(consumerName); await sut.ActivateAsync(); @@ -214,7 +214,7 @@ namespace Squidex.Infrastructure.EventSourcing.Grains [Fact] public async Task Should_not_invoke_and_update_position_when_event_is_from_another_subscription() { - var @event = new StoredEvent(Guid.NewGuid().ToString(), 123, eventData); + var @event = new StoredEvent("Stream", Guid.NewGuid().ToString(), 123, eventData); await sut.OnActivateAsync(consumerName); await sut.ActivateAsync(); @@ -302,7 +302,7 @@ namespace Squidex.Infrastructure.EventSourcing.Grains A.CallTo(() => eventConsumer.On(envelope)) .Throws(ex); - var @event = new StoredEvent(Guid.NewGuid().ToString(), 123, eventData); + var @event = new StoredEvent("Stream", Guid.NewGuid().ToString(), 123, eventData); await sut.OnActivateAsync(consumerName); await sut.ActivateAsync(); @@ -329,7 +329,7 @@ namespace Squidex.Infrastructure.EventSourcing.Grains A.CallTo(() => formatter.Parse(eventData, true)) .Throws(ex); - var @event = new StoredEvent(Guid.NewGuid().ToString(), 123, eventData); + var @event = new StoredEvent("Stream", Guid.NewGuid().ToString(), 123, eventData); await sut.OnActivateAsync(consumerName); await sut.ActivateAsync(); @@ -356,7 +356,7 @@ namespace Squidex.Infrastructure.EventSourcing.Grains A.CallTo(() => eventConsumer.On(envelope)) .Throws(exception); - var @event = new StoredEvent(Guid.NewGuid().ToString(), 123, eventData); + var @event = new StoredEvent("Stream", Guid.NewGuid().ToString(), 123, eventData); await sut.OnActivateAsync(consumerName); await sut.ActivateAsync(); diff --git a/tests/Squidex.Infrastructure.Tests/EventSourcing/RetrySubscriptionTests.cs b/tests/Squidex.Infrastructure.Tests/EventSourcing/RetrySubscriptionTests.cs index 74a7eca02..85aee50d2 100644 --- a/tests/Squidex.Infrastructure.Tests/EventSourcing/RetrySubscriptionTests.cs +++ b/tests/Squidex.Infrastructure.Tests/EventSourcing/RetrySubscriptionTests.cs @@ -90,7 +90,7 @@ namespace Squidex.Infrastructure.EventSourcing [Fact] public async Task Should_forward_event_from_inner_subscription() { - var ev = new StoredEvent("1", 2, new EventData()); + var ev = new StoredEvent("Stream", "1", 2, new EventData()); await OnEventAsync(eventSubscription, ev); await sut.StopAsync(); @@ -102,7 +102,7 @@ namespace Squidex.Infrastructure.EventSourcing [Fact] public async Task Should_not_forward_event_when_message_is_from_another_subscription() { - var ev = new StoredEvent("1", 2, new EventData()); + var ev = new StoredEvent("Stream", "1", 2, new EventData()); await OnEventAsync(A.Fake(), ev); await sut.StopAsync(); diff --git a/tests/Squidex.Infrastructure.Tests/States/PersistenceEventSourcingTests.cs b/tests/Squidex.Infrastructure.Tests/States/PersistenceEventSourcingTests.cs index 3b23fb421..5957ed378 100644 --- a/tests/Squidex.Infrastructure.Tests/States/PersistenceEventSourcingTests.cs +++ b/tests/Squidex.Infrastructure.Tests/States/PersistenceEventSourcingTests.cs @@ -56,7 +56,7 @@ namespace Squidex.Infrastructure.States [Fact] public async Task Should_ignore_old_events() { - var storedEvent = new StoredEvent("1", 0, new EventData()); + var storedEvent = new StoredEvent("1", "1", 0, new EventData()); A.CallTo(() => eventStore.QueryAsync(key, 0)) .Returns(new List { storedEvent }); @@ -252,7 +252,7 @@ namespace Squidex.Infrastructure.States foreach (var @event in events) { var eventData = new EventData(); - var eventStored = new StoredEvent(i.ToString(), i, eventData); + var eventStored = new StoredEvent(i.ToString(), i.ToString(), i, eventData); eventsStored.Add(eventStored);