Browse Source

Single handlers to be responsible for all backup, restore and cleanup belongings.

pull/311/head
Sebastian 8 years ago
parent
commit
7c1c88d4a2
  1. 5
      src/Squidex.Domain.Apps.Entities.MongoDb/Assets/MongoAssetRepository.cs
  2. 5
      src/Squidex.Domain.Apps.Entities.MongoDb/Contents/MongoContentPublishedCollection.cs
  3. 7
      src/Squidex.Domain.Apps.Entities.MongoDb/Contents/MongoContentRepository.cs
  4. 5
      src/Squidex.Domain.Apps.Entities.MongoDb/History/MongoHistoryEventRepository.cs
  5. 36
      src/Squidex.Domain.Apps.Entities/AppGrainCleaner.cs
  6. 3
      src/Squidex.Domain.Apps.Entities/AppProvider.cs
  7. 8
      src/Squidex.Domain.Apps.Entities/Apps/AppGrain.cs
  8. 91
      src/Squidex.Domain.Apps.Entities/Apps/BackupApps.cs
  9. 9
      src/Squidex.Domain.Apps.Entities/Apps/Guards/GuardApp.cs
  10. 40
      src/Squidex.Domain.Apps.Entities/Apps/Indexes/AppsByNameIndexCommandMiddleware.cs
  11. 28
      src/Squidex.Domain.Apps.Entities/Apps/Indexes/AppsByNameIndexGrain.cs
  12. 6
      src/Squidex.Domain.Apps.Entities/Apps/Indexes/IAppsByNameIndex.cs
  13. 2
      src/Squidex.Domain.Apps.Entities/Apps/Indexes/IAppsByUserIndex.cs
  14. 124
      src/Squidex.Domain.Apps.Entities/Assets/BackupAssets.cs
  15. 2
      src/Squidex.Domain.Apps.Entities/Assets/Repositories/IAssetRepository.cs
  16. 138
      src/Squidex.Domain.Apps.Entities/Backup/AppCleanerGrain.cs
  17. 53
      src/Squidex.Domain.Apps.Entities/Backup/BackupGrain.cs
  18. 57
      src/Squidex.Domain.Apps.Entities/Backup/BackupHandler.cs
  19. 11
      src/Squidex.Domain.Apps.Entities/Backup/BackupHandlerWithStore.cs
  20. 70
      src/Squidex.Domain.Apps.Entities/Backup/BackupReader.cs
  21. 59
      src/Squidex.Domain.Apps.Entities/Backup/BackupWriter.cs
  22. 13
      src/Squidex.Domain.Apps.Entities/Backup/CleanerStatus.cs
  23. 50
      src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreApp.cs
  24. 72
      src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreAssets.cs
  25. 2
      src/Squidex.Domain.Apps.Entities/Backup/IAppCleanerGrain.cs
  26. 24
      src/Squidex.Domain.Apps.Entities/Backup/IRestoreHandler.cs
  27. 191
      src/Squidex.Domain.Apps.Entities/Backup/RestoreGrain.cs
  28. 31
      src/Squidex.Domain.Apps.Entities/Contents/BackupContents.cs
  29. 2
      src/Squidex.Domain.Apps.Entities/Contents/Repositories/IContentRepository.cs
  30. 32
      src/Squidex.Domain.Apps.Entities/History/BackupHistory.cs
  31. 2
      src/Squidex.Domain.Apps.Entities/History/Repositories/IHistoryEventRepository.cs
  32. 40
      src/Squidex.Domain.Apps.Entities/Rules/BackupRules.cs
  33. 2
      src/Squidex.Domain.Apps.Entities/Rules/Indexes/IRulesByAppIndex.cs
  34. 40
      src/Squidex.Domain.Apps.Entities/Schemas/BackupSchemas.cs
  35. 2
      src/Squidex.Domain.Apps.Entities/Schemas/Indexes/ISchemasByAppIndex.cs
  36. 6
      src/Squidex.Domain.Apps.Entities/Tags/GrainTagService.cs
  37. 2
      src/Squidex.Domain.Apps.Entities/Tags/ITagService.cs
  38. 5
      src/Squidex.Infrastructure/Log/Profiler.cs
  39. 11
      src/Squidex/Config/Domain/EntitiesServices.cs
  40. 46
      tests/Squidex.Domain.Apps.Entities.Tests/AppGrainCleanerTests.cs
  41. 6
      tests/Squidex.Domain.Apps.Entities.Tests/Apps/AppGrainTests.cs
  42. 27
      tests/Squidex.Domain.Apps.Entities.Tests/Apps/Guards/GuardAppTests.cs
  43. 32
      tests/Squidex.Domain.Apps.Entities.Tests/Apps/Indexes/AppsByNameIndexCommandMiddlewareTests.cs
  44. 64
      tests/Squidex.Domain.Apps.Entities.Tests/Backup/EventStreamTests.cs
  45. 2
      tests/Squidex.Domain.Apps.Entities.Tests/Tags/GrainTagServiceTests.cs
  46. 3
      tools/Migrate_01/Migrations/PopulateGrainIndexes.cs

5
src/Squidex.Domain.Apps.Entities.MongoDb/Assets/MongoAssetRepository.cs

@ -120,5 +120,10 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Assets
return assetEntity;
}
}
public Task RemoveAsync(Guid appId)
{
return Collection.DeleteManyAsync(x => x.IndexedAppId == appId);
}
}
}

5
src/Squidex.Domain.Apps.Entities.MongoDb/Contents/MongoContentPublishedCollection.cs

@ -56,10 +56,5 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Contents
return Collection.ReplaceOneAsync(x => x.Id == content.Id, content, new UpdateOptions { IsUpsert = true });
}
public Task RemoveAsync(Guid id)
{
return Collection.DeleteOneAsync(x => x.Id == id);
}
}
}

7
src/Squidex.Domain.Apps.Entities.MongoDb/Contents/MongoContentRepository.cs

@ -115,6 +115,13 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.Contents
}
}
public Task RemoveAsync(Guid appId)
{
return Task.WhenAll(
contentsDraft.RemoveAsync(appId),
contentsPublished.RemoveAsync(appId));
}
public Task ClearAsync()
{
return Task.WhenAll(

5
src/Squidex.Domain.Apps.Entities.MongoDb/History/MongoHistoryEventRepository.cs

@ -112,5 +112,10 @@ namespace Squidex.Domain.Apps.Entities.MongoDb.History
}
}
}
public Task RemoveAsync(Guid appId)
{
return Collection.DeleteManyAsync(x => x.AppId == appId);
}
}
}

36
src/Squidex.Domain.Apps.Entities/AppGrainCleaner.cs

@ -1,36 +0,0 @@
// ==========================================================================
// Squidex Headless CMS
// ==========================================================================
// Copyright (c) Squidex UG (haftungsbeschraenkt)
// All rights reserved. Licensed under the MIT license.
// ==========================================================================
using System;
using System.Threading.Tasks;
using Orleans;
using Squidex.Infrastructure;
namespace Squidex.Domain.Apps.Entities
{
public sealed class AppGrainCleaner<T> : ICleanableAppStorage where T : ICleanableAppGrain
{
private readonly IGrainFactory grainFactory;
public string Name
{
get { return typeof(T).Name; }
}
public AppGrainCleaner(IGrainFactory grainFactory)
{
Guard.NotNull(grainFactory, nameof(grainFactory));
this.grainFactory = grainFactory;
}
public Task ClearAsync(Guid appId)
{
return grainFactory.GetGrain<T>(appId).ClearAsync();
}
}
}

3
src/Squidex.Domain.Apps.Entities/AppProvider.cs

@ -11,8 +11,11 @@ using System.Linq;
using System.Threading.Tasks;
using Orleans;
using Squidex.Domain.Apps.Entities.Apps;
using Squidex.Domain.Apps.Entities.Apps.Indexes;
using Squidex.Domain.Apps.Entities.Rules;
using Squidex.Domain.Apps.Entities.Rules.Indexes;
using Squidex.Domain.Apps.Entities.Schemas;
using Squidex.Domain.Apps.Entities.Schemas.Indexes;
using Squidex.Infrastructure;
using Squidex.Infrastructure.Caching;
using Squidex.Infrastructure.Log;

8
src/Squidex.Domain.Apps.Entities/Apps/AppGrain.cs

@ -29,7 +29,6 @@ namespace Squidex.Domain.Apps.Entities.Apps
public sealed class AppGrain : SquidexDomainObjectGrain<AppState>, IAppGrain
{
private readonly InitialPatterns initialPatterns;
private readonly IAppProvider appProvider;
private readonly IAppPlansProvider appPlansProvider;
private readonly IAppPlanBillingManager appPlansBillingManager;
private readonly IUserResolver userResolver;
@ -38,20 +37,17 @@ namespace Squidex.Domain.Apps.Entities.Apps
InitialPatterns initialPatterns,
IStore<Guid> store,
ISemanticLog log,
IAppProvider appProvider,
IAppPlansProvider appPlansProvider,
IAppPlanBillingManager appPlansBillingManager,
IUserResolver userResolver)
: base(store, log)
{
Guard.NotNull(initialPatterns, nameof(initialPatterns));
Guard.NotNull(appProvider, nameof(appProvider));
Guard.NotNull(userResolver, nameof(userResolver));
Guard.NotNull(appPlansProvider, nameof(appPlansProvider));
Guard.NotNull(appPlansBillingManager, nameof(appPlansBillingManager));
this.userResolver = userResolver;
this.appProvider = appProvider;
this.appPlansProvider = appPlansProvider;
this.appPlansBillingManager = appPlansBillingManager;
this.initialPatterns = initialPatterns;
@ -64,9 +60,9 @@ namespace Squidex.Domain.Apps.Entities.Apps
switch (command)
{
case CreateApp createApp:
return CreateAsync(createApp, async c =>
return CreateAsync(createApp, c =>
{
await GuardApp.CanCreate(c, appProvider);
GuardApp.CanCreate(c);
Create(c);
});

91
src/Squidex.Domain.Apps.Entities/Apps/BackupApps.cs

@ -0,0 +1,91 @@
// ==========================================================================
// 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 Orleans;
using Squidex.Domain.Apps.Entities.Apps.Indexes;
using Squidex.Domain.Apps.Entities.Apps.State;
using Squidex.Domain.Apps.Entities.Backup;
using Squidex.Domain.Apps.Events.Apps;
using Squidex.Infrastructure;
using Squidex.Infrastructure.EventSourcing;
using Squidex.Infrastructure.Orleans;
using Squidex.Infrastructure.States;
namespace Squidex.Domain.Apps.Entities.Apps
{
public sealed class BackupApps : BackupHandlerWithStore
{
private readonly IGrainFactory grainFactory;
private readonly HashSet<string> users = new HashSet<string>();
private bool isReserved;
private AppCreated appCreated;
public BackupApps(IStore<Guid> store, IGrainFactory grainFactory)
: base(store)
{
Guard.NotNull(grainFactory, nameof(grainFactory));
this.grainFactory = grainFactory;
}
public override Task RemoveAsync(Guid appId)
{
return RemoveSnapshotAsync<AppState>(appId);
}
public async override Task RestoreEventAsync(Envelope<IEvent> @event, Guid appId, BackupReader reader)
{
switch (@event.Payload)
{
case AppCreated appCreated:
{
this.appCreated = appCreated;
var index = grainFactory.GetGrain<IAppsByNameIndex>(SingleGrain.Id);
if (!(isReserved = await index.ReserveAppAsync(appCreated.AppId.Id, appCreated.AppId.Name)))
{
throw new DomainException("The app id or name is not available.");
}
break;
}
case AppContributorAssigned contributorAssigned:
users.Add(contributorAssigned.ContributorId);
break;
case AppContributorRemoved contributorRemoved:
users.Remove(contributorRemoved.ContributorId);
break;
}
}
public override async Task RestoreAsync(Guid appId, BackupReader reader)
{
await grainFactory.GetGrain<IAppsByNameIndex>(SingleGrain.Id).AddAppAsync(appCreated.AppId.Id, appCreated.AppId.Name);
foreach (var user in users)
{
await grainFactory.GetGrain<IAppsByUserIndex>(user).AddAppAsync(appCreated.AppId.Id);
}
}
public override async Task CleanupRestoreAsync(Guid appId, Exception exception)
{
if (isReserved)
{
var index = grainFactory.GetGrain<IAppsByNameIndex>(SingleGrain.Id);
await index.ReserveAppAsync(appCreated.AppId.Id, appCreated.AppId.Name);
}
}
}
}

9
src/Squidex.Domain.Apps.Entities/Apps/Guards/GuardApp.cs

@ -6,7 +6,6 @@
// ==========================================================================
using System;
using System.Threading.Tasks;
using Squidex.Domain.Apps.Core.Apps;
using Squidex.Domain.Apps.Entities.Apps.Commands;
using Squidex.Domain.Apps.Entities.Apps.Services;
@ -16,20 +15,16 @@ namespace Squidex.Domain.Apps.Entities.Apps.Guards
{
public static class GuardApp
{
public static Task CanCreate(CreateApp command, IAppProvider appProvider)
public static void CanCreate(CreateApp command)
{
Guard.NotNull(command, nameof(command));
return Validate.It(() => "Cannot create app.", async e =>
Validate.It(() => "Cannot create app.", e =>
{
if (!command.Name.IsSlug())
{
e("Name must be a valid slug.", nameof(command.Name));
}
else if (await appProvider.GetAppAsync(command.Name) != null)
{
e("An app with the same name already exists.", nameof(command.Name));
}
});
}

40
src/Squidex.Domain.Apps.Entities/Apps/Indexes/AppsByNameIndexCommandMiddleware.cs

@ -28,20 +28,44 @@ namespace Squidex.Domain.Apps.Entities.Apps.Indexes
public async Task HandleAsync(CommandContext context, Func<Task> next)
{
if (context.IsCompleted)
var createApp = context.Command as CreateApp;
var isReserved = false;
try
{
switch (context.Command)
if (createApp != null)
{
isReserved = await index.ReserveAppAsync(createApp.AppId, createApp.Name);
if (!isReserved)
{
var error = new ValidationError("An app with the same name already exists.", nameof(createApp.Name));
throw new ValidationException("Cannot create app.", error);
}
}
await next();
if (context.IsCompleted)
{
case CreateApp createApp:
if (createApp != null)
{
await index.AddAppAsync(createApp.AppId, createApp.Name);
break;
case ArchiveApp archiveApp:
}
else if (context.Command is ArchiveApp archiveApp)
{
await index.RemoveAppAsync(archiveApp.AppId);
break;
}
}
}
finally
{
if (isReserved && createApp != null)
{
await index.RemoveReservationAsync(createApp.AppId, createApp.Name);
}
}
await next();
}
}
}

28
src/Squidex.Domain.Apps.Entities/Apps/Indexes/AppsByNameIndexGrain.cs

@ -12,12 +12,15 @@ using System.Threading.Tasks;
using Squidex.Infrastructure;
using Squidex.Infrastructure.Orleans;
using Squidex.Infrastructure.States;
using Squidex.Infrastructure.Tasks;
namespace Squidex.Domain.Apps.Entities.Apps.Indexes
{
public sealed class AppsByNameIndexGrain : GrainOfString, IAppsByNameIndex
{
private readonly IStore<string> store;
private readonly HashSet<Guid> reservedIds = new HashSet<Guid>();
private readonly HashSet<string> reservedNames = new HashSet<string>();
private IPersistence<State> persistence;
private State state = new State();
@ -51,6 +54,31 @@ namespace Squidex.Domain.Apps.Entities.Apps.Indexes
return persistence.WriteSnapshotAsync(state);
}
public Task<bool> ReserveAppAsync(Guid appId, string name)
{
var canReserve =
!state.Apps.ContainsKey(name) &&
!state.Apps.Any(x => x.Value == appId) &&
!reservedIds.Contains(appId) &&
!reservedNames.Contains(name);
if (canReserve)
{
reservedIds.Add(appId);
reservedNames.Add(name);
}
return Task.FromResult(canReserve);
}
public Task RemoveReservationAsync(Guid appId, string name)
{
reservedIds.Remove(appId);
reservedNames.Remove(name);
return TaskHelper.Done;
}
public Task AddAppAsync(Guid appId, string name)
{
state.Apps[name] = appId;

6
src/Squidex.Domain.Apps.Entities/Apps/Indexes/IAppsByNameIndex.cs

@ -10,16 +10,20 @@ using System.Collections.Generic;
using System.Threading.Tasks;
using Orleans;
namespace Squidex.Domain.Apps.Entities.Apps
namespace Squidex.Domain.Apps.Entities.Apps.Indexes
{
public interface IAppsByNameIndex : IGrainWithStringKey
{
Task<bool> ReserveAppAsync(Guid appId, string name);
Task AddAppAsync(Guid appId, string name);
Task RemoveAppAsync(Guid appId);
Task RebuildAsync(Dictionary<string, Guid> apps);
Task RemoveReservationAsync(Guid appId, string name);
Task<Guid> GetAppIdAsync(string name);
Task<List<Guid>> GetAppIdsAsync();

2
src/Squidex.Domain.Apps.Entities/Apps/Indexes/IAppsByUserIndex.cs

@ -10,7 +10,7 @@ using System.Collections.Generic;
using System.Threading.Tasks;
using Orleans;
namespace Squidex.Domain.Apps.Entities.Apps
namespace Squidex.Domain.Apps.Entities.Apps.Indexes
{
public interface IAppsByUserIndex : IGrainWithStringKey
{

124
src/Squidex.Domain.Apps.Entities/Assets/BackupAssets.cs

@ -0,0 +1,124 @@
// ==========================================================================
// 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.Domain.Apps.Entities.Assets.Repositories;
using Squidex.Domain.Apps.Entities.Assets.State;
using Squidex.Domain.Apps.Entities.Backup;
using Squidex.Domain.Apps.Entities.Tags;
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.Assets
{
public sealed class BackupAssets : BackupHandlerWithStore
{
private readonly HashSet<Guid> assetIds = new HashSet<Guid>();
private readonly IAssetStore assetStore;
private readonly IAssetRepository assetRepository;
private readonly ITagService tagService;
private readonly IEventDataFormatter eventDataFormatter;
public BackupAssets(IStore<Guid> store,
IEventDataFormatter eventDataFormatter,
IAssetStore assetStore,
IAssetRepository assetRepository,
ITagService tagService)
: base(store)
{
Guard.NotNull(eventDataFormatter, nameof(eventDataFormatter));
Guard.NotNull(assetStore, nameof(assetStore));
Guard.NotNull(assetRepository, nameof(assetRepository));
Guard.NotNull(tagService, nameof(tagService));
this.eventDataFormatter = eventDataFormatter;
this.assetStore = assetStore;
this.assetRepository = assetRepository;
this.tagService = tagService;
}
public override async Task RemoveAsync(Guid appId)
{
await tagService.ClearAsync(appId, TagGroups.Assets);
await assetRepository.RemoveAsync(appId);
}
public override Task BackupEventAsync(EventData @event, Guid appId, BackupWriter writer)
{
if (@event.Type == "AssetCreatedEvent" ||
@event.Type == "AssetUpdatedEvent")
{
var parsedEvent = eventDataFormatter.Parse(@event);
switch (parsedEvent.Payload)
{
case AssetCreated assetCreated:
return WriteAssetAsync(assetCreated.AssetId, assetCreated.FileVersion, writer);
case AssetUpdated assetUpdated:
return WriteAssetAsync(assetUpdated.AssetId, assetUpdated.FileVersion, writer);
}
}
return TaskHelper.Done;
}
public override Task RestoreEventAsync(Envelope<IEvent> @event, Guid appId, BackupReader reader)
{
switch (@event.Payload)
{
case AssetCreated assetCreated:
assetIds.Add(assetCreated.AssetId);
return ReadAssetAsync(assetCreated.AssetId, assetCreated.FileVersion, reader);
case AssetUpdated assetUpdated:
return ReadAssetAsync(assetUpdated.AssetId, assetUpdated.FileVersion, reader);
}
return TaskHelper.Done;
}
public override Task RestoreAsync(Guid appId, BackupReader reader)
{
return RebuildManyAsync(assetIds, id => RebuildAsync<AssetState, AssetGrain>(id, (e, s) => s.Apply(e)));
}
private Task WriteAssetAsync(Guid assetId, long fileVersion, BackupWriter writer)
{
return writer.WriteAttachmentAsync(GetName(assetId, fileVersion), stream =>
{
return assetStore.DownloadAsync(assetId.ToString(), fileVersion, null, stream);
});
}
private Task ReadAssetAsync(Guid assetId, long fileVersion, BackupReader reader)
{
return reader.ReadAttachmentAsync(GetName(assetId, fileVersion), async stream =>
{
try
{
await assetStore.UploadAsync(assetId.ToString(), fileVersion, null, stream);
}
catch (AssetAlreadyExistsException)
{
return;
}
});
}
private static string GetName(Guid assetId, long fileVersion)
{
return $"{assetId}_{fileVersion}.asset";
}
}
}

2
src/Squidex.Domain.Apps.Entities/Assets/Repositories/IAssetRepository.cs

@ -19,5 +19,7 @@ namespace Squidex.Domain.Apps.Entities.Assets.Repositories
Task<IResultList<IAssetEntity>> QueryAsync(Guid appId, HashSet<Guid> ids);
Task<IAssetEntity> FindAssetAsync(Guid id);
Task RemoveAsync(Guid appId);
}
}

138
src/Squidex.Domain.Apps.Entities/Backup/AppCleanerGrain.cs

@ -12,11 +12,6 @@ using System.Threading.Tasks;
using Orleans;
using Orleans.Concurrency;
using Orleans.Runtime;
using Squidex.Domain.Apps.Entities.Apps.State;
using Squidex.Domain.Apps.Entities.Rules;
using Squidex.Domain.Apps.Entities.Rules.State;
using Squidex.Domain.Apps.Entities.Schemas;
using Squidex.Domain.Apps.Entities.Schemas.State;
using Squidex.Infrastructure;
using Squidex.Infrastructure.EventSourcing;
using Squidex.Infrastructure.Log;
@ -32,7 +27,7 @@ namespace Squidex.Domain.Apps.Entities.Backup
private readonly IGrainFactory grainFactory;
private readonly IStore<Guid> store;
private readonly IEventStore eventStore;
private readonly IEnumerable<ICleanableAppStorage> storages;
private readonly IEnumerable<BackupHandler> handlers;
private readonly ISemanticLog log;
private IPersistence<State> persistence;
private bool isCleaning;
@ -43,21 +38,21 @@ namespace Squidex.Domain.Apps.Entities.Backup
{
public HashSet<Guid> Apps { get; set; } = new HashSet<Guid>();
public HashSet<Guid> PendingApps { get; set; } = new HashSet<Guid>();
public HashSet<Guid> FailedApps { get; set; } = new HashSet<Guid>();
}
public AppCleanerGrain(IGrainFactory grainFactory, IEventStore eventStore, IStore<Guid> store, IEnumerable<ICleanableAppStorage> storages, ISemanticLog log)
public AppCleanerGrain(IGrainFactory grainFactory, IEventStore eventStore, IStore<Guid> store, IEnumerable<BackupHandler> handlers, ISemanticLog log)
{
Guard.NotNull(grainFactory, nameof(grainFactory));
Guard.NotNull(store, nameof(store));
Guard.NotNull(storages, nameof(storages));
Guard.NotNull(handlers, nameof(handlers));
Guard.NotNull(eventStore, nameof(eventStore));
Guard.NotNull(log, nameof(log));
this.grainFactory = grainFactory;
this.store = store;
this.storages = storages;
this.handlers = handlers;
this.log = log;
@ -99,6 +94,22 @@ namespace Squidex.Domain.Apps.Entities.Backup
return TaskHelper.Done;
}
public Task<CleanerStatus> GetStatusAsync(Guid appId)
{
if (state.Apps.Contains(appId))
{
return Task.FromResult(CleanerStatus.Cleaning);
}
else if (state.FailedApps.Contains(appId))
{
return Task.FromResult(CleanerStatus.Failed);
}
else
{
return Task.FromResult(CleanerStatus.Cleaned);
}
}
private async Task CleanAsync()
{
if (isCleaning)
@ -111,43 +122,7 @@ namespace Squidex.Domain.Apps.Entities.Backup
{
foreach (var appId in state.Apps.ToList())
{
using (Profiler.StartSession())
{
try
{
log.LogInformation(w => w
.WriteProperty("action", "cleanApp")
.WriteProperty("status", "started")
.WriteProperty("appId", appId.ToString()));
await CleanAsync(appId);
state.Apps.Remove(appId);
log.LogInformation(w =>
{
w.WriteProperty("action", "cleanApp");
w.WriteProperty("status", "completed");
w.WriteProperty("appId", appId.ToString());
Profiler.Session?.Write(w);
});
}
catch (Exception ex)
{
state.PendingApps.Add(appId);
log.LogError(ex, w => w
.WriteProperty("action", "cleanApp")
.WriteProperty("appId", appId.ToString()));
}
finally
{
state.Apps.Remove(appId);
await persistence.WriteSnapshotAsync(state);
}
}
await CleanupAppAsync(appId);
}
}
finally
@ -156,47 +131,64 @@ namespace Squidex.Domain.Apps.Entities.Backup
}
}
private async Task CleanAsync(Guid appId)
private async Task CleanupAppAsync(Guid appId)
{
using (Profiler.Trace("DeleteEvents"))
using (Profiler.StartSession())
{
await eventStore.DeleteManyAsync("AppId", appId);
}
try
{
log.LogInformation(w => w
.WriteProperty("action", "cleanApp")
.WriteProperty("status", "started")
.WriteProperty("appId", appId.ToString()));
using (Profiler.Trace("DeleteRules"))
{
var ruleIds = await grainFactory.GetGrain<IRulesByAppIndex>(appId).GetRuleIdsAsync();
await CleanupCoreAsync(appId);
foreach (var ruleId in ruleIds)
{
await store.RemoveSnapshotAsync<RuleState>(ruleId);
log.LogInformation(w =>
{
w.WriteProperty("action", "cleanApp");
w.WriteProperty("status", "completed");
w.WriteProperty("appId", appId.ToString());
Profiler.Session?.Write(w);
});
}
}
catch (Exception ex)
{
state.FailedApps.Add(appId);
using (Profiler.Trace("DeleteSchemas"))
{
var schemaIds = await grainFactory.GetGrain<ISchemasByAppIndex>(appId).GetSchemaIdsAsync();
log.LogError(ex, w =>
{
w.WriteProperty("action", "cleanApp");
w.WriteProperty("status", "failed");
w.WriteProperty("appId", appId.ToString());
foreach (var schemaId in schemaIds)
Profiler.Session?.Write(w);
});
}
finally
{
await store.RemoveSnapshotAsync<SchemaState>(schemaId);
state.Apps.Remove(appId);
await persistence.WriteSnapshotAsync(state);
}
}
}
foreach (var storage in storages)
private async Task CleanupCoreAsync(Guid appId)
{
using (Profiler.Trace("DeleteEvents"))
{
using (Profiler.Trace($"{storage.Name}.ClearAsync"))
await eventStore.DeleteManyAsync("AppId", appId);
}
foreach (var handler in handlers)
{
using (Profiler.TraceMethod(handler.GetType(), nameof(BackupHandler.RemoveAsync)))
{
await storage.ClearAsync(appId);
await handler.RemoveAsync(appId);
}
}
await store.RemoveSnapshotAsync<AppState>(appId);
}
private async Task DeleteAsync<TState>(Guid id)
{
await store.RemoveSnapshotAsync<TState>(id);
}
}
}

53
src/Squidex.Domain.Apps.Entities/Backup/BackupGrain.cs

@ -14,7 +14,6 @@ using NodaTime;
using Orleans.Concurrency;
using Squidex.Domain.Apps.Entities.Backup.State;
using Squidex.Domain.Apps.Events;
using Squidex.Domain.Apps.Events.Assets;
using Squidex.Infrastructure;
using Squidex.Infrastructure.Assets;
using Squidex.Infrastructure.EventSourcing;
@ -32,6 +31,7 @@ namespace Squidex.Domain.Apps.Entities.Backup
private readonly IClock clock;
private readonly IAssetStore assetStore;
private readonly IEventDataFormatter eventDataFormatter;
private readonly IEnumerable<BackupHandler> handlers;
private readonly ISemanticLog log;
private readonly IEventStore eventStore;
private readonly IBackupArchiveLocation backupArchiveLocation;
@ -48,6 +48,7 @@ namespace Squidex.Domain.Apps.Entities.Backup
IClock clock,
IEventStore eventStore,
IEventDataFormatter eventDataFormatter,
IEnumerable<BackupHandler> handlers,
ISemanticLog log,
IStore<Guid> store)
{
@ -56,6 +57,7 @@ namespace Squidex.Domain.Apps.Entities.Backup
Guard.NotNull(clock, nameof(clock));
Guard.NotNull(eventStore, nameof(eventStore));
Guard.NotNull(eventDataFormatter, nameof(eventDataFormatter));
Guard.NotNull(handlers, nameof(handlers));
Guard.NotNull(store, nameof(store));
Guard.NotNull(log, nameof(log));
@ -64,6 +66,7 @@ namespace Squidex.Domain.Apps.Entities.Backup
this.clock = clock;
this.eventStore = eventStore;
this.eventDataFormatter = eventDataFormatter;
this.handlers = handlers;
this.store = store;
this.log = log;
}
@ -176,45 +179,21 @@ namespace Squidex.Domain.Apps.Entities.Backup
{
using (var stream = await backupArchiveLocation.OpenStreamAsync(job.Id))
{
using (var writer = new EventStreamWriter(stream))
using (var writer = new BackupWriter(stream))
{
await eventStore.QueryAsync(async @event =>
{
var eventData = @event.Data;
if (eventData.Type == "AssetCreatedEvent" ||
eventData.Type == "AssetUpdatedEvent")
{
var parsedEvent = eventDataFormatter.Parse(eventData);
var assetVersion = 0L;
var assetId = Guid.Empty;
switch (parsedEvent.Payload)
{
case AssetCreated assetCreated:
assetId = assetCreated.AssetId;
assetVersion = assetCreated.FileVersion;
break;
case AssetUpdated asetUpdated:
assetId = asetUpdated.AssetId;
assetVersion = asetUpdated.FileVersion;
break;
}
await writer.WriteEventAsync(@event, attachment =>
{
return assetStore.DownloadAsync(assetId.ToString(), assetVersion, null, attachment);
});
job.HandledAssets++;
}
else
writer.WriteEvent(@event);
foreach (var handler in handlers)
{
await writer.WriteEventAsync(@event);
await handler.BackupEventAsync(eventData, appId, writer);
}
job.HandledEvents++;
job.HandledEvents = writer.WrittenEvents;
job.HandledAssets = writer.WrittenAttachments;
var now = clock.GetCurrentInstant();
@ -225,6 +204,16 @@ namespace Squidex.Domain.Apps.Entities.Backup
await WriteAsync();
}
}, SquidexHeaders.AppId, appId.ToString(), null, currentTask.Token);
foreach (var handler in handlers)
{
await handler.BackupAsync(appId, writer);
}
foreach (var handler in handlers)
{
await handler.CompleteBackupAsync(appId, writer);
}
}
stream.Position = 0;

57
src/Squidex.Domain.Apps.Entities/Backup/BackupHandler.cs

@ -0,0 +1,57 @@
// ==========================================================================
// 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.EventSourcing;
using Squidex.Infrastructure.Tasks;
namespace Squidex.Domain.Apps.Entities.Backup
{
public abstract class BackupHandler
{
public virtual Task RestoreEventAsync(Envelope<IEvent> @event, Guid appId, BackupReader reader)
{
return TaskHelper.Done;
}
public virtual Task BackupEventAsync(EventData @event, Guid appId, BackupWriter writer)
{
return TaskHelper.Done;
}
public virtual Task RestoreAsync(Guid appId, BackupReader reader)
{
return TaskHelper.Done;
}
public virtual Task BackupAsync(Guid appId, BackupWriter writer)
{
return TaskHelper.Done;
}
public virtual Task CleanupRestoreAsync(Guid appId, Exception exception)
{
return TaskHelper.Done;
}
public virtual Task CompleteRestoreAsync(Guid appId, BackupReader reader)
{
return TaskHelper.Done;
}
public virtual Task CompleteBackupAsync(Guid appId, BackupWriter writer)
{
return TaskHelper.Done;
}
public virtual Task RemoveAsync(Guid appId)
{
return TaskHelper.Done;
}
}
}

11
src/Squidex.Domain.Apps.Entities/Backup/Handlers/HandlerBase.cs → src/Squidex.Domain.Apps.Entities/Backup/BackupHandlerWithStore.cs

@ -13,19 +13,24 @@ using Squidex.Infrastructure.Commands;
using Squidex.Infrastructure.EventSourcing;
using Squidex.Infrastructure.States;
namespace Squidex.Domain.Apps.Entities.Backup.Handlers
namespace Squidex.Domain.Apps.Entities.Backup
{
public abstract class HandlerBase
public abstract class BackupHandlerWithStore : BackupHandler
{
private readonly IStore<Guid> store;
protected HandlerBase(IStore<Guid> store)
protected BackupHandlerWithStore(IStore<Guid> store)
{
Guard.NotNull(store, nameof(store));
this.store = store;
}
protected Task RemoveSnapshotAsync<TState>(Guid id)
{
return store.RemoveSnapshotAsync<TState>(id);
}
protected async Task RebuildManyAsync(IEnumerable<Guid> ids, Func<Guid, Task> action)
{
foreach (var id in ids)

70
src/Squidex.Domain.Apps.Entities/Backup/EventStreamReader.cs → src/Squidex.Domain.Apps.Entities/Backup/BackupReader.cs

@ -15,13 +15,26 @@ using Squidex.Infrastructure.EventSourcing;
namespace Squidex.Domain.Apps.Entities.Backup
{
public sealed class EventStreamReader : DisposableObjectBase
public sealed class BackupReader : DisposableObjectBase
{
private const int MaxItemsPerFolder = 1000;
private const int MaxEventsPerFolder = 1000;
private const int MaxAttachmentFolders = 1000;
private static readonly JsonSerializer JsonSerializer = JsonSerializer.CreateDefault();
private readonly ZipArchive archive;
private int readEvents;
private int readAttachments;
public EventStreamReader(Stream stream)
public int ReadEvents
{
get { return readEvents; }
}
public int ReadAttachments
{
get { return readAttachments; }
}
public BackupReader(Stream stream)
{
archive = new ZipArchive(stream, ZipArchiveMode.Read, true);
}
@ -34,16 +47,35 @@ namespace Squidex.Domain.Apps.Entities.Backup
}
}
public async Task ReadEventsAsync(Func<StoredEvent, Stream, Task> eventHandler)
public async Task ReadAttachmentAsync(string name, Func<Stream, Task> handler)
{
Guard.NotNull(eventHandler, nameof(eventHandler));
Guard.NotNullOrEmpty(name, nameof(name));
Guard.NotNull(handler, nameof(handler));
var readEvents = 0;
var readAttachments = 0;
var attachmentFolder = Math.Abs(name.GetHashCode() % MaxAttachmentFolders);
var attachmentPath = $"attachments/{attachmentFolder}/{name}";
var attachmentEntry = archive.GetEntry(attachmentPath);
if (attachmentEntry == null)
{
throw new FileNotFoundException("Cannot find attachment.", name);
}
using (var stream = attachmentEntry.Open())
{
await handler(stream);
}
readAttachments++;
}
public async Task ReadEventsAsync(Func<StoredEvent, Task> handler)
{
Guard.NotNull(handler, nameof(handler));
while (true)
{
var eventFolder = readEvents / MaxItemsPerFolder;
var eventFolder = readEvents / MaxEventsPerFolder;
var eventPath = $"events/{eventFolder}/{readEvents}.json";
var eventEntry = archive.GetEntry(eventPath);
@ -52,33 +84,15 @@ namespace Squidex.Domain.Apps.Entities.Backup
break;
}
StoredEvent eventData;
using (var stream = eventEntry.Open())
{
using (var textReader = new StreamReader(stream))
{
eventData = (StoredEvent)JsonSerializer.Deserialize(textReader, typeof(StoredEvent));
}
}
var attachmentFolder = readAttachments / MaxItemsPerFolder;
var attachmentPath = $"attachments/{attachmentFolder}/{readEvents}.blob";
var attachmentEntry = archive.GetEntry(attachmentPath);
if (attachmentEntry != null)
{
using (var stream = attachmentEntry.Open())
{
await eventHandler(eventData, stream);
var storedEvent = (StoredEvent)JsonSerializer.Deserialize(textReader, typeof(StoredEvent));
readAttachments++;
await handler(storedEvent);
}
}
else
{
await eventHandler(eventData, null);
}
readEvents++;
}

59
src/Squidex.Domain.Apps.Entities/Backup/EventStreamWriter.cs → src/Squidex.Domain.Apps.Entities/Backup/BackupWriter.cs

@ -10,20 +10,31 @@ using System.IO;
using System.IO.Compression;
using System.Threading.Tasks;
using Newtonsoft.Json;
using Newtonsoft.Json.Linq;
using Squidex.Infrastructure;
using Squidex.Infrastructure.EventSourcing;
namespace Squidex.Domain.Apps.Entities.Backup
{
public sealed class EventStreamWriter : DisposableObjectBase
public sealed class BackupWriter : DisposableObjectBase
{
private const int MaxItemsPerFolder = 1000;
private const int MaxEventsPerFolder = 1000;
private const int MaxAttachmentFolders = 1000;
private static readonly JsonSerializer JsonSerializer = JsonSerializer.CreateDefault();
private readonly ZipArchive archive;
private int writtenEvents;
private int writtenAttachments;
public EventStreamWriter(Stream stream)
public int WrittenEvents
{
get { return writtenEvents; }
}
public int WrittenAttachments
{
get { return writtenAttachments; }
}
public BackupWriter(Stream stream)
{
archive = new ZipArchive(stream, ZipArchiveMode.Update, true);
}
@ -36,11 +47,26 @@ namespace Squidex.Domain.Apps.Entities.Backup
}
}
public async Task WriteEventAsync(StoredEvent storedEvent, Func<Stream, Task> attachment = null)
public async Task WriteAttachmentAsync(string name, Func<Stream, Task> handler)
{
var eventObject = JObject.FromObject(storedEvent);
Guard.NotNullOrEmpty(name, nameof(name));
Guard.NotNull(handler, nameof(handler));
var attachmentFolder = Math.Abs(name.GetHashCode() % MaxAttachmentFolders);
var attachmentPath = $"attachments/{attachmentFolder}/{name}";
var attachmentEntry = archive.CreateEntry(attachmentPath);
using (var stream = attachmentEntry.Open())
{
await handler(stream);
}
var eventFolder = writtenEvents / MaxItemsPerFolder;
writtenAttachments++;
}
public void WriteEvent(StoredEvent storedEvent)
{
var eventFolder = writtenEvents / MaxEventsPerFolder;
var eventPath = $"events/{eventFolder}/{writtenEvents}.json";
var eventEntry = archive.GetEntry(eventPath) ?? archive.CreateEntry(eventPath);
@ -48,25 +74,8 @@ namespace Squidex.Domain.Apps.Entities.Backup
{
using (var textWriter = new StreamWriter(stream))
{
using (var jsonWriter = new JsonTextWriter(textWriter))
{
await eventObject.WriteToAsync(jsonWriter);
}
}
}
if (attachment != null)
{
var attachmentFolder = writtenAttachments / MaxItemsPerFolder;
var attachmentPath = $"attachments/{attachmentFolder}/{writtenEvents}.blob";
var attachmentEntry = archive.GetEntry(attachmentPath) ?? archive.CreateEntry(attachmentPath);
using (var stream = attachmentEntry.Open())
{
await attachment(stream);
JsonSerializer.Serialize(textWriter, storedEvent);
}
writtenAttachments++;
}
writtenEvents++;

13
src/Squidex.Domain.Apps.Entities/ICleanableAppStorage.cs → src/Squidex.Domain.Apps.Entities/Backup/CleanerStatus.cs

@ -5,15 +5,12 @@
// All rights reserved. Licensed under the MIT license.
// ==========================================================================
using System;
using System.Threading.Tasks;
namespace Squidex.Domain.Apps.Entities
namespace Squidex.Domain.Apps.Entities.Backup
{
public interface ICleanableAppStorage
public enum CleanerStatus
{
string Name { get; }
Task ClearAsync(Guid appId);
Cleaned,
Cleaning,
Failed
}
}

50
src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreApp.cs

@ -1,50 +0,0 @@
// ==========================================================================
// 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<Guid> appId;
public string Name { get; } = "App";
public RestoreApp(IStore<Guid> store)
: base(store)
{
}
public Task HandleAsync(Envelope<IEvent> @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;
}
}
}

72
src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreAssets.cs

@ -1,72 +0,0 @@
// ==========================================================================
// 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<Guid> assetIds = new HashSet<Guid>();
private readonly IAssetStore assetStore;
public string Name { get; } = "Assets";
public RestoreAssets(IStore<Guid> store, IAssetStore assetStore)
: base(store)
{
Guard.NotNull(assetStore, nameof(assetStore));
this.assetStore = assetStore;
}
public async Task HandleAsync(Envelope<IEvent> @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<AssetState, AssetGrain>(id, (e, s) => s.Apply(e)));
}
public Task CompleteAsync()
{
return TaskHelper.Done;
}
}
}

2
src/Squidex.Domain.Apps.Entities/Backup/IAppCleanerGrain.cs

@ -14,5 +14,7 @@ namespace Squidex.Domain.Apps.Entities.Backup
public interface IAppCleanerGrain : IBackgroundGrain
{
Task EnqueueAppAsync(Guid appId);
Task<CleanerStatus> GetStatusAsync(Guid appId);
}
}

24
src/Squidex.Domain.Apps.Entities/Backup/IRestoreHandler.cs

@ -1,24 +0,0 @@
// ==========================================================================
// 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<IEvent> @event, Stream attachment);
Task ProcessAsync();
Task CompleteAsync();
}
}

191
src/Squidex.Domain.Apps.Entities/Backup/RestoreGrain.cs

@ -13,6 +13,7 @@ using NodaTime;
using Orleans;
using Squidex.Domain.Apps.Entities.Backup.State;
using Squidex.Domain.Apps.Events;
using Squidex.Domain.Apps.Events.Apps;
using Squidex.Infrastructure;
using Squidex.Infrastructure.Assets;
using Squidex.Infrastructure.EventSourcing;
@ -29,12 +30,13 @@ namespace Squidex.Domain.Apps.Entities.Backup
private readonly IClock clock;
private readonly IAssetStore assetStore;
private readonly IEventDataFormatter eventDataFormatter;
private readonly IGrainFactory grainFactory;
private readonly IAppCleanerGrain appCleaner;
private readonly ISemanticLog log;
private readonly IEventStore eventStore;
private readonly IBackupArchiveLocation backupArchiveLocation;
private readonly IStore<string> store;
private readonly IEnumerable<IRestoreHandler> handlers;
private readonly IEnumerable<BackupHandler> handlers;
private RestoreState state = new RestoreState();
private IPersistence<RestoreState> persistence;
@ -45,9 +47,9 @@ namespace Squidex.Domain.Apps.Entities.Backup
IEventStore eventStore,
IEventDataFormatter eventDataFormatter,
IGrainFactory grainFactory,
IEnumerable<IRestoreHandler> handlers,
IEnumerable<BackupHandler> handlers,
ISemanticLog log,
IStore<Guid> store)
IStore<string> store)
{
Guard.NotNull(assetStore, nameof(assetStore));
Guard.NotNull(backupArchiveLocation, nameof(backupArchiveLocation));
@ -64,9 +66,12 @@ namespace Squidex.Domain.Apps.Entities.Backup
this.clock = clock;
this.eventStore = eventStore;
this.eventDataFormatter = eventDataFormatter;
this.grainFactory = grainFactory;
this.handlers = handlers;
this.store = store;
this.log = log;
appCleaner = grainFactory.GetGrain<IAppCleanerGrain>(SingleGrain.Id);
}
public override async Task OnActivateAsync(string key)
@ -97,10 +102,7 @@ namespace Squidex.Domain.Apps.Entities.Backup
state.Job.Status = "Failed due application restart";
state.Job.IsFailed = true;
if (state.Job.AppId != Guid.Empty)
{
appCleaner.EnqueueAppAsync(state.Job.AppId).Forget();
}
TryCleanup();
await persistence.WriteSnapshotAsync(state);
}
@ -108,48 +110,96 @@ namespace Squidex.Domain.Apps.Entities.Backup
private async Task ProcessAsync()
{
try
using (Profiler.StartSession())
{
await DoAsync(
"Downloading Backup",
"Downloaded Backup",
DownloadAsync);
try
{
log.LogInformation(w => w
.WriteProperty("action", "restore")
.WriteProperty("status", "started")
.WriteProperty("url", state.Job.Uri.ToString()));
await DoAsync(
"Reading Events",
"Readed Events",
ReadEventsAsync);
state.Job.Status = "Downloading Backup";
foreach (var handler in handlers)
{
await DoAsync($"{handler.Name} Proessing", $"{handler.Name} Processed", handler.ProcessAsync);
}
using (Profiler.Trace("Download"))
{
await DownloadAsync();
}
foreach (var handler in handlers)
{
await DoAsync($"{handler.Name} Completing", $"{handler.Name} Completed", handler.CompleteAsync);
state.Job.Status = "Downloaded Backup";
using (var stream = await backupArchiveLocation.OpenStreamAsync(state.Job.Id))
{
using (var reader = new BackupReader(stream))
{
using (Profiler.Trace("ReadEvents"))
{
await ReadEventsAsync(reader);
}
state.Job.Status = "Events read";
foreach (var handler in handlers)
{
using (Profiler.TraceMethod(handler.GetType(), nameof(BackupHandler.RestoreAsync)))
{
await handler.RestoreAsync(state.Job.AppId, reader);
}
state.Job.Status = $"{handler} Processed";
}
foreach (var handler in handlers)
{
using (Profiler.TraceMethod(handler.GetType(), nameof(BackupHandler.CompleteRestoreAsync)))
{
await handler.CompleteRestoreAsync(state.Job.AppId, reader);
}
state.Job.Status = $"{handler} Completed";
}
}
}
state.Job = null;
log.LogInformation(w =>
{
w.WriteProperty("action", "restore");
w.WriteProperty("status", "completed");
w.WriteProperty("url", state.Job.Uri.ToString());
Profiler.Session?.Write(w);
});
}
catch (Exception ex)
{
state.Job.IsFailed = true;
state.Job = null;
}
catch (Exception ex)
{
log.LogError(ex, w => w
.WriteProperty("action", "makeBackup")
.WriteProperty("status", "failed")
.WriteProperty("backupId", state.Job.Id.ToString()));
if (state.Job.AppId != Guid.Empty)
{
foreach (var handler in handlers)
{
await handler.CleanupRestoreAsync(state.Job.AppId, ex);
}
}
state.Job.IsFailed = true;
TryCleanup();
log.LogError(ex, w =>
{
w.WriteProperty("action", "retore");
w.WriteProperty("status", "failed");
w.WriteProperty("url", state.Job.Uri.ToString());
if (state.Job.AppId != Guid.Empty)
Profiler.Session?.Write(w);
});
}
finally
{
appCleaner.EnqueueAppAsync(state.Job.AppId).Forget();
await persistence.WriteSnapshotAsync(state);
}
}
finally
{
await persistence.WriteSnapshotAsync(state);
}
}
private async Task DownloadAsync()
@ -166,47 +216,58 @@ namespace Squidex.Domain.Apps.Entities.Backup
}
}
private async Task ReadEventsAsync()
private async Task ReadEventsAsync(BackupReader reader)
{
using (var stream = await backupArchiveLocation.OpenStreamAsync(state.Job.Id))
await reader.ReadEventsAsync(async (storedEvent) =>
{
using (var reader = new EventStreamReader(stream))
var eventData = storedEvent.Data;
var eventParsed = eventDataFormatter.Parse(eventData);
if (eventParsed.Payload is SquidexEvent squidexEvent)
{
squidexEvent.Actor = state.Job.User;
}
else if (eventParsed.Payload is AppCreated appCreated)
{
var eventIndex = 0;
state.Job.AppId = appCreated.AppId.Id;
await reader.ReadEventsAsync(async (@event, attachment) =>
{
var eventData = @event.Data;
await CheckCleanupStatus();
}
var parsedEvent = eventDataFormatter.Parse(eventData);
foreach (var handler in handlers)
{
await handler.RestoreEventAsync(eventParsed, state.Job.AppId, reader);
}
if (parsedEvent.Payload is SquidexEvent squidexEvent)
{
squidexEvent.Actor = state.Job.User;
}
await eventStore.AppendAsync(Guid.NewGuid(), storedEvent.StreamName, new List<EventData> { storedEvent.Data });
foreach (var handler in handlers)
{
await handler.HandleAsync(parsedEvent, attachment);
}
state.Job.Status = $"Handled event {reader.ReadEvents} events and {reader.ReadAttachments} attachments";
});
}
await eventStore.AppendAsync(Guid.NewGuid(), @event.StreamName, new List<EventData> { @event.Data });
private async Task CheckCleanupStatus()
{
var cleaner = grainFactory.GetGrain<IAppCleanerGrain>(SingleGrain.Id);
eventIndex++;
var status = await cleaner.GetStatusAsync(state.Job.AppId);
state.Job.Status = $"Handled event {eventIndex}";
});
}
if (status == CleanerStatus.Cleaning)
{
throw new DomainException("The app is removed in the background.");
}
if (status == CleanerStatus.Cleaning)
{
throw new DomainException("The app could not be cleaned.");
}
}
private async Task DoAsync(string start, string end, Func<Task> action)
private void TryCleanup()
{
state.Job.Status = start;
await action();
state.Job.Status = end;
if (state.Job.AppId != Guid.Empty)
{
appCleaner.EnqueueAppAsync(state.Job.AppId).Forget();
}
}
public Task<J<IRestoreJob>> GetStateAsync()

31
src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreContents.cs → src/Squidex.Domain.Apps.Entities/Contents/BackupContents.cs

@ -7,29 +7,37 @@
using System;
using System.Collections.Generic;
using System.IO;
using System.Threading.Tasks;
using Squidex.Domain.Apps.Entities.Contents;
using Squidex.Domain.Apps.Entities.Backup;
using Squidex.Domain.Apps.Entities.Contents.Repositories;
using Squidex.Domain.Apps.Entities.Contents.State;
using Squidex.Domain.Apps.Events.Contents;
using Squidex.Infrastructure;
using Squidex.Infrastructure.EventSourcing;
using Squidex.Infrastructure.States;
using Squidex.Infrastructure.Tasks;
namespace Squidex.Domain.Apps.Entities.Backup.Handlers
namespace Squidex.Domain.Apps.Entities.Contents
{
public sealed class RestoreContents : HandlerBase, IRestoreHandler
public sealed class BackupContents : BackupHandlerWithStore
{
private readonly HashSet<Guid> contentIds = new HashSet<Guid>();
private readonly IContentRepository contentRepository;
public string Name { get; } = "Contents";
public RestoreContents(IStore<Guid> store)
public BackupContents(IStore<Guid> store, IContentRepository contentRepository)
: base(store)
{
Guard.NotNull(contentRepository, nameof(contentRepository));
this.contentRepository = contentRepository;
}
public override Task RemoveAsync(Guid appId)
{
return contentRepository.RemoveAsync(appId);
}
public Task HandleAsync(Envelope<IEvent> @event, Stream attachment)
public override Task RestoreEventAsync(Envelope<IEvent> @event, Guid appId, BackupReader reader)
{
switch (@event.Payload)
{
@ -41,14 +49,9 @@ namespace Squidex.Domain.Apps.Entities.Backup.Handlers
return TaskHelper.Done;
}
public Task ProcessAsync()
public override Task RestoreAsync(Guid appId, BackupReader reader)
{
return RebuildManyAsync(contentIds, id => RebuildAsync<ContentState, ContentGrain>(id, (e, s) => s.Apply(e)));
}
public Task CompleteAsync()
{
return TaskHelper.Done;
}
}
}

2
src/Squidex.Domain.Apps.Entities/Contents/Repositories/IContentRepository.cs

@ -28,5 +28,7 @@ namespace Squidex.Domain.Apps.Entities.Contents.Repositories
Task<IContentEntity> FindContentAsync(IAppEntity app, ISchemaEntity schema, Status[] status, Guid id);
Task QueryScheduledWithoutDataAsync(Instant now, Func<IContentEntity, Task> callback);
Task RemoveAsync(Guid appId);
}
}

32
src/Squidex.Domain.Apps.Entities/History/BackupHistory.cs

@ -0,0 +1,32 @@
// ==========================================================================
// Squidex Headless CMS
// ==========================================================================
// Copyright (c) Squidex UG (haftungsbeschraenkt)
// All rights reserved. Licensed under the MIT license.
// ==========================================================================
using System;
using System.Threading.Tasks;
using Squidex.Domain.Apps.Entities.Backup;
using Squidex.Domain.Apps.Entities.History.Repositories;
using Squidex.Infrastructure;
namespace Squidex.Domain.Apps.Entities.History
{
public sealed class BackupHistory : BackupHandler
{
private readonly IHistoryEventRepository historyEventRepository;
public BackupHistory(IHistoryEventRepository historyEventRepository)
{
Guard.NotNull(historyEventRepository, nameof(historyEventRepository));
this.historyEventRepository = historyEventRepository;
}
public override Task RemoveAsync(Guid appId)
{
return historyEventRepository.RemoveAsync(appId);
}
}
}

2
src/Squidex.Domain.Apps.Entities/History/Repositories/IHistoryEventRepository.cs

@ -14,5 +14,7 @@ namespace Squidex.Domain.Apps.Entities.History.Repositories
public interface IHistoryEventRepository
{
Task<IReadOnlyList<IHistoryEventEntity>> QueryByChannelAsync(Guid appId, string channelPrefix, int count);
Task RemoveAsync(Guid appId);
}
}

40
src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreRules.cs → src/Squidex.Domain.Apps.Entities/Rules/BackupRules.cs

@ -7,29 +7,25 @@
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.Backup;
using Squidex.Domain.Apps.Entities.Rules.Indexes;
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
namespace Squidex.Domain.Apps.Entities.Rules
{
public sealed class RestoreRules : HandlerBase, IRestoreHandler
public sealed class BackupRules : BackupHandlerWithStore
{
private readonly HashSet<Guid> ruleIds = new HashSet<Guid>();
private readonly IGrainFactory grainFactory;
private Guid appId;
public string Name { get; } = "Rules";
public RestoreRules(IStore<Guid> store, IGrainFactory grainFactory)
public BackupRules(IStore<Guid> store, IGrainFactory grainFactory)
: base(store)
{
Guard.NotNull(grainFactory, nameof(grainFactory));
@ -37,13 +33,24 @@ namespace Squidex.Domain.Apps.Entities.Backup.Handlers
this.grainFactory = grainFactory;
}
public Task HandleAsync(Envelope<IEvent> @event, Stream attachment)
public override async Task RemoveAsync(Guid appId)
{
var index = grainFactory.GetGrain<IRulesByAppIndex>(appId);
var idsToRemove = await grainFactory.GetGrain<IRulesByAppIndex>(appId).GetRuleIdsAsync();
foreach (var ruleId in idsToRemove)
{
await RemoveSnapshotAsync<RuleState>(ruleId);
}
await index.ClearAsync();
}
public override Task RestoreEventAsync(Envelope<IEvent> @event, Guid appId, BackupReader reader)
{
switch (@event.Payload)
{
case AppCreated appCreated:
appId = appCreated.AppId.Id;
break;
case RuleCreated ruleCreated:
ruleIds.Add(ruleCreated.RuleId);
break;
@ -55,16 +62,11 @@ namespace Squidex.Domain.Apps.Entities.Backup.Handlers
return TaskHelper.Done;
}
public async Task ProcessAsync()
public async override Task RestoreAsync(Guid appId, BackupReader reader)
{
await RebuildManyAsync(ruleIds, id => RebuildAsync<RuleState, RuleGrain>(id, (e, s) => s.Apply(e)));
await grainFactory.GetGrain<IRulesByAppIndex>(appId).RebuildAsync(ruleIds);
}
public Task CompleteAsync()
{
return TaskHelper.Done;
}
}
}

2
src/Squidex.Domain.Apps.Entities/Rules/Indexes/IRulesByAppIndex.cs

@ -9,7 +9,7 @@ using System;
using System.Collections.Generic;
using System.Threading.Tasks;
namespace Squidex.Domain.Apps.Entities.Rules
namespace Squidex.Domain.Apps.Entities.Rules.Indexes
{
public interface IRulesByAppIndex : ICleanableAppGrain
{

40
src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreSchemas.cs → src/Squidex.Domain.Apps.Entities/Schemas/BackupSchemas.cs

@ -7,33 +7,29 @@
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.Backup;
using Squidex.Domain.Apps.Entities.Schemas.Indexes;
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
namespace Squidex.Domain.Apps.Entities.Schemas
{
public sealed class RestoreSchemas : HandlerBase, IRestoreHandler
public sealed class BackupSchemas : BackupHandlerWithStore
{
private readonly HashSet<NamedId<Guid>> schemaIds = new HashSet<NamedId<Guid>>();
private readonly Dictionary<string, Guid> schemasByName = new Dictionary<string, Guid>();
private readonly FieldRegistry fieldRegistry;
private readonly IGrainFactory grainFactory;
private Guid appId;
public string Name { get; } = "Schemas";
public RestoreSchemas(IStore<Guid> store, FieldRegistry fieldRegistry, IGrainFactory grainFactory)
public BackupSchemas(IStore<Guid> store, FieldRegistry fieldRegistry, IGrainFactory grainFactory)
: base(store)
{
Guard.NotNull(fieldRegistry, nameof(fieldRegistry));
@ -44,13 +40,24 @@ namespace Squidex.Domain.Apps.Entities.Backup.Handlers
this.grainFactory = grainFactory;
}
public Task HandleAsync(Envelope<IEvent> @event, Stream attachment)
public override async Task RemoveAsync(Guid appId)
{
var index = grainFactory.GetGrain<ISchemasByAppIndex>(appId);
var idsToRemove = await index.GetSchemaIdsAsync();
foreach (var schemaId in idsToRemove)
{
await RemoveSnapshotAsync<SchemaState>(schemaId);
}
await index.ClearAsync();
}
public override Task RestoreEventAsync(Envelope<IEvent> @event, Guid appId, BackupReader reader)
{
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;
@ -60,16 +67,11 @@ namespace Squidex.Domain.Apps.Entities.Backup.Handlers
return TaskHelper.Done;
}
public async Task ProcessAsync()
public async override Task RestoreAsync(Guid appId, BackupReader reader)
{
await RebuildManyAsync(schemaIds.Select(x => x.Id), id => RebuildAsync<SchemaState, SchemaGrain>(id, (e, s) => s.Apply(e, fieldRegistry)));
await grainFactory.GetGrain<ISchemasByAppIndex>(appId).RebuildAsync(schemasByName);
}
public Task CompleteAsync()
{
return TaskHelper.Done;
}
}
}

2
src/Squidex.Domain.Apps.Entities/Schemas/Indexes/ISchemasByAppIndex.cs

@ -9,7 +9,7 @@ using System;
using System.Collections.Generic;
using System.Threading.Tasks;
namespace Squidex.Domain.Apps.Entities.Schemas
namespace Squidex.Domain.Apps.Entities.Schemas.Indexes
{
public interface ISchemasByAppIndex : ICleanableAppGrain
{

6
src/Squidex.Domain.Apps.Entities/Tags/GrainTagService.cs

@ -13,7 +13,7 @@ using Squidex.Infrastructure;
namespace Squidex.Domain.Apps.Entities.Tags
{
public sealed class GrainTagService : ITagService, ICleanableAppStorage
public sealed class GrainTagService : ITagService
{
private readonly IGrainFactory grainFactory;
@ -54,9 +54,9 @@ namespace Squidex.Domain.Apps.Entities.Tags
return GetGrain(appId, group).RebuildTagsAsync(allTags);
}
public Task ClearAsync(Guid appId)
public Task ClearAsync(Guid appId, string group)
{
return GetGrain(appId, TagGroups.Assets).ClearAsync();
return GetGrain(appId, group).ClearAsync();
}
private ITagGrain GetGrain(Guid appId, string group)

2
src/Squidex.Domain.Apps.Entities/Tags/ITagService.cs

@ -22,5 +22,7 @@ namespace Squidex.Domain.Apps.Entities.Tags
Task<Dictionary<string, int>> GetTagsAsync(Guid appId, string group);
Task RebuildTagsAsync(Guid appId, string group, Dictionary<string, string> allTags);
Task ClearAsync(Guid appId, string group);
}
}

5
src/Squidex.Infrastructure/Log/Profiler.cs

@ -34,6 +34,11 @@ namespace Squidex.Infrastructure.Log
return Cleaner;
}
public static IDisposable TraceMethod(Type type, [CallerMemberName] string memberName = null)
{
return Trace($"{type.Name}/{memberName}");
}
public static IDisposable TraceMethod<T>([CallerMemberName] string memberName = null)
{
return Trace($"{typeof(T).Name}/{memberName}");

11
src/Squidex/Config/Domain/EntitiesServices.cs

@ -85,7 +85,7 @@ namespace Squidex.Config.Domain
.AsSelf();
services.AddSingletonAs<GrainTagService>()
.As<ITagService>().As<ICleanableAppStorage>();
.As<ITagService>();
services.AddSingletonAs<FileTypeTagGenerator>()
.As<ITagGenerator<CreateAsset>>();
@ -93,15 +93,6 @@ namespace Squidex.Config.Domain
services.AddSingletonAs<ImageTagGenerator>()
.As<ITagGenerator<CreateAsset>>();
services.AddSingletonAs<AppGrainCleaner<IBackupGrain>>()
.As<ICleanableAppStorage>();
services.AddSingletonAs<AppGrainCleaner<IRulesByAppIndex>>()
.As<ICleanableAppStorage>();
services.AddSingletonAs<AppGrainCleaner<ISchemasByAppIndex>>()
.As<ICleanableAppStorage>();
services.AddSingletonAs<JintScriptEngine>()
.As<IScriptEngine>();

46
tests/Squidex.Domain.Apps.Entities.Tests/AppGrainCleanerTests.cs

@ -1,46 +0,0 @@
// ==========================================================================
// Squidex Headless CMS
// ==========================================================================
// Copyright (c) Squidex UG (haftungsbeschraenkt)
// All rights reserved. Licensed under the MIT license.
// ==========================================================================
using System;
using System.Threading.Tasks;
using FakeItEasy;
using Orleans;
using Xunit;
namespace Squidex.Domain.Apps.Entities
{
public class AppGrainCleanerTests
{
private readonly IGrainFactory grainFactory = A.Fake<IGrainFactory>();
private readonly ICleanableAppGrain index = A.Fake<ICleanableAppGrain>();
private readonly Guid appId = Guid.NewGuid();
private readonly AppGrainCleaner<ICleanableAppGrain> sut;
public AppGrainCleanerTests()
{
A.CallTo(() => grainFactory.GetGrain<ICleanableAppGrain>(appId, null))
.Returns(index);
sut = new AppGrainCleaner<ICleanableAppGrain>(grainFactory);
}
[Fact]
public void Should_provide_name()
{
Assert.Equal(typeof(ICleanableAppGrain).Name, sut.Name);
}
[Fact]
public async Task Should_forward_to_index()
{
await sut.ClearAsync(appId);
A.CallTo(() => index.ClearAsync())
.MustHaveHappened();
}
}
}

6
tests/Squidex.Domain.Apps.Entities.Tests/Apps/AppGrainTests.cs

@ -25,7 +25,6 @@ namespace Squidex.Domain.Apps.Entities.Apps
{
public class AppGrainTests : HandlerTestBase<AppGrain, AppState>
{
private readonly IAppProvider appProvider = A.Fake<IAppProvider>();
private readonly IAppPlansProvider appPlansProvider = A.Fake<IAppPlansProvider>();
private readonly IAppPlanBillingManager appPlansBillingManager = A.Fake<IAppPlanBillingManager>();
private readonly IUser user = A.Fake<IUser>();
@ -47,9 +46,6 @@ namespace Squidex.Domain.Apps.Entities.Apps
public AppGrainTests()
{
A.CallTo(() => appProvider.GetAppAsync(AppName))
.Returns((IAppEntity)null);
A.CallTo(() => user.Id)
.Returns(contributorId);
@ -62,7 +58,7 @@ namespace Squidex.Domain.Apps.Entities.Apps
{ patternId2, new AppPattern("Numbers", "[0-9]*") }
};
sut = new AppGrain(initialPatterns, Store, A.Dummy<ISemanticLog>(), appProvider, appPlansProvider, appPlansBillingManager, userResolver);
sut = new AppGrain(initialPatterns, Store, A.Dummy<ISemanticLog>(), appPlansProvider, appPlansBillingManager, userResolver);
sut.OnActivateAsync(Id).Wait();
}

27
tests/Squidex.Domain.Apps.Entities.Tests/Apps/Guards/GuardAppTests.cs

@ -5,8 +5,8 @@
// All rights reserved. Licensed under the MIT license.
// ==========================================================================
using System.Threading.Tasks;
using FakeItEasy;
using Orleans;
using Squidex.Domain.Apps.Core.Apps;
using Squidex.Domain.Apps.Entities.Apps.Commands;
using Squidex.Domain.Apps.Entities.Apps.Services;
@ -19,18 +19,12 @@ namespace Squidex.Domain.Apps.Entities.Apps.Guards
{
public class GuardAppTests
{
private readonly IAppProvider apps = A.Fake<IAppProvider>();
private readonly IUserResolver users = A.Fake<IUserResolver>();
private readonly IGrainFactory grainFactory = A.Fake<IGrainFactory>();
private readonly IAppPlansProvider appPlans = A.Fake<IAppPlansProvider>();
public GuardAppTests()
{
A.CallTo(() => apps.GetAppAsync(A<string>.Ignored))
.Returns(Task.FromResult<IAppEntity>(null));
A.CallTo(() => apps.GetAppAsync("existing"))
.Returns(A.Dummy<IAppEntity>());
A.CallTo(() => users.FindByIdOrEmailAsync(A<string>.Ignored))
.Returns(A.Dummy<IUser>());
@ -42,29 +36,20 @@ namespace Squidex.Domain.Apps.Entities.Apps.Guards
}
[Fact]
public Task CanCreate_should_throw_exception_if_name_already_in_use()
{
var command = new CreateApp { Name = "existing" };
return ValidationAssert.ThrowsAsync(() => GuardApp.CanCreate(command, apps),
new ValidationError("An app with the same name already exists.", "Name"));
}
[Fact]
public Task CanCreate_should_throw_exception_if_name_not_valid()
public void CanCreate_should_throw_exception_if_name_not_valid()
{
var command = new CreateApp { Name = "INVALID NAME" };
return ValidationAssert.ThrowsAsync(() => GuardApp.CanCreate(command, apps),
ValidationAssert.Throws(() => GuardApp.CanCreate(command),
new ValidationError("Name must be a valid slug.", "Name"));
}
[Fact]
public Task CanCreate_should_not_throw_exception_if_app_name_is_free()
public void CanCreate_should_not_throw_exception_if_app_name_is_valid()
{
var command = new CreateApp { Name = "new-app" };
return GuardApp.CanCreate(command, apps);
GuardApp.CanCreate(command);
}
[Fact]

32
tests/Squidex.Domain.Apps.Entities.Tests/Apps/Indexes/AppsByNameIndexCommandMiddlewareTests.cs

@ -10,6 +10,7 @@ using System.Threading.Tasks;
using FakeItEasy;
using Orleans;
using Squidex.Domain.Apps.Entities.Apps.Commands;
using Squidex.Infrastructure;
using Squidex.Infrastructure.Commands;
using Squidex.Infrastructure.Orleans;
using Xunit;
@ -35,14 +36,45 @@ namespace Squidex.Domain.Apps.Entities.Apps.Indexes
[Fact]
public async Task Should_add_app_to_index_on_create()
{
A.CallTo(() => index.ReserveAppAsync(appId, "my-app"))
.Returns(true);
var context =
new CommandContext(new CreateApp { AppId = appId, Name = "my-app" }, commandBus)
.Complete();
await sut.HandleAsync(context);
A.CallTo(() => index.ReserveAppAsync(appId, "my-app"))
.MustHaveHappened();
A.CallTo(() => index.AddAppAsync(appId, "my-app"))
.MustHaveHappened();
A.CallTo(() => index.RemoveReservationAsync(appId, "my-app"))
.MustHaveHappened();
}
[Fact]
public async Task Should_not_remove_reservation_when_not_reserved()
{
A.CallTo(() => index.ReserveAppAsync(appId, "my-app"))
.Returns(false);
var context =
new CommandContext(new CreateApp { AppId = appId, Name = "my-app" }, commandBus)
.Complete();
await Assert.ThrowsAsync<ValidationException>(() => sut.HandleAsync(context));
A.CallTo(() => index.ReserveAppAsync(appId, "my-app"))
.MustHaveHappened();
A.CallTo(() => index.AddAppAsync(appId, "my-app"))
.MustNotHaveHappened();
A.CallTo(() => index.RemoveReservationAsync(appId, "my-app"))
.MustNotHaveHappened();
}
[Fact]

64
tests/Squidex.Domain.Apps.Entities.Tests/Backup/EventStreamTests.cs

@ -10,72 +10,64 @@ using System.IO;
using System.Threading.Tasks;
using FluentAssertions;
using Squidex.Infrastructure.EventSourcing;
using Squidex.Infrastructure.Tasks;
using Xunit;
namespace Squidex.Domain.Apps.Entities.Backup
{
public class EventStreamTests
{
public sealed class EventInfo
{
public StoredEvent Stored { get; set; }
public byte[] Attachment { get; set; }
}
[Fact]
public async Task Should_write_and_read_events()
{
var stream = new MemoryStream();
var sourceEvents = new List<EventInfo>();
var sourceEvents = new List<StoredEvent>();
for (var i = 0; i < 1000; i++)
using (var writer = new BackupWriter(stream))
{
var eventData = new EventData { Type = i.ToString(), Metadata = i, Payload = i };
var eventInfo = new EventInfo { Stored = new StoredEvent("S", "1", 2, eventData) };
if (i % 10 == 0)
for (var i = 0; i < 1000; i++)
{
eventInfo.Attachment = new byte[] { (byte)i };
}
sourceEvents.Add(eventInfo);
}
var eventData = new EventData { Type = i.ToString(), Metadata = i, Payload = i };
var eventStored = new StoredEvent("S", "1", 2, eventData);
using (var reader = new EventStreamWriter(stream))
{
foreach (var @event in sourceEvents)
{
if (@event.Attachment == null)
{
await reader.WriteEventAsync(@event.Stored);
}
else
if (i % 10 == 0)
{
await reader.WriteEventAsync(@event.Stored, s => s.WriteAsync(@event.Attachment, 0, 1));
await writer.WriteAttachmentAsync(eventData.Type, innerStream =>
{
return innerStream.WriteAsync(new byte[] { (byte)i }, 0, 1);
});
}
writer.WriteEvent(eventStored);
sourceEvents.Add(eventStored);
}
}
stream.Position = 0;
var readEvents = new List<EventInfo>();
var readEvents = new List<StoredEvent>();
using (var reader = new EventStreamReader(stream))
using (var reader = new BackupReader(stream))
{
await reader.ReadEventsAsync(async (stored, attachment) =>
await reader.ReadEventsAsync(async @event =>
{
var eventInfo = new EventInfo { Stored = stored };
var i = int.Parse(@event.Data.Type);
if (attachment != null)
if (i % 10 == 0)
{
eventInfo.Attachment = new byte[1];
await reader.ReadAttachmentAsync(@event.Data.Type, innerStream =>
{
var b = innerStream.ReadByte();
Assert.Equal((byte)i, b);
await attachment.ReadAsync(eventInfo.Attachment, 0, 1);
return TaskHelper.Done;
});
}
readEvents.Add(eventInfo);
readEvents.Add(@event);
});
}

2
tests/Squidex.Domain.Apps.Entities.Tests/Tags/GrainTagServiceTests.cs

@ -38,7 +38,7 @@ namespace Squidex.Domain.Apps.Entities.Tags
[Fact]
public async Task Should_call_grain_when_clearing()
{
await sut.ClearAsync(appId);
await sut.ClearAsync(appId, TagGroups.Assets);
A.CallTo(() => grain.ClearAsync())
.MustHaveHappened();

3
tools/Migrate_01/Migrations/PopulateGrainIndexes.cs

@ -9,11 +9,8 @@ using System;
using System.Collections.Generic;
using System.Threading.Tasks;
using Orleans;
using Squidex.Domain.Apps.Entities.Apps;
using Squidex.Domain.Apps.Entities.Apps.State;
using Squidex.Domain.Apps.Entities.Rules;
using Squidex.Domain.Apps.Entities.Rules.State;
using Squidex.Domain.Apps.Entities.Schemas;
using Squidex.Domain.Apps.Entities.Schemas.State;
using Squidex.Infrastructure;
using Squidex.Infrastructure.Migrations;

Loading…
Cancel
Save