Browse Source

Extracted restore behavior.

pull/311/head
Sebastian 8 years ago
parent
commit
191b29c1c7
  1. 25
      src/Squidex.Domain.Apps.Entities/Backup/BackupGrain.cs
  2. 10
      src/Squidex.Domain.Apps.Entities/Backup/EventStreamReader.cs
  3. 12
      src/Squidex.Domain.Apps.Entities/Backup/EventStreamWriter.cs
  4. 55
      src/Squidex.Domain.Apps.Entities/Backup/Handlers/HandlerBase.cs
  5. 50
      src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreApp.cs
  6. 72
      src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreAssets.cs
  7. 54
      src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreContents.cs
  8. 70
      src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreRules.cs
  9. 75
      src/Squidex.Domain.Apps.Entities/Backup/Handlers/RestoreSchemas.cs
  10. 21
      src/Squidex.Domain.Apps.Entities/Backup/IRestoreGrain.cs
  11. 24
      src/Squidex.Domain.Apps.Entities/Backup/IRestoreHandler.cs
  12. 23
      src/Squidex.Domain.Apps.Entities/Backup/IRestoreJob.cs
  13. 203
      src/Squidex.Domain.Apps.Entities/Backup/RestoreGrain.cs
  14. 17
      src/Squidex.Domain.Apps.Entities/Backup/State/RestoreState.cs
  15. 38
      src/Squidex.Domain.Apps.Entities/Backup/State/RestoreStateJob.cs
  16. 5
      src/Squidex.Domain.Apps.Entities/Tags/GrainTagService.cs
  17. 1
      src/Squidex.Domain.Apps.Entities/Tags/ITagGrain.cs
  18. 2
      src/Squidex.Domain.Apps.Entities/Tags/ITagService.cs
  19. 5
      src/Squidex.Domain.Apps.Entities/Tags/TagGrain.cs
  20. 1
      src/Squidex.Infrastructure.GetEventStore/EventSourcing/Formatter.cs
  21. 4
      src/Squidex.Infrastructure.MongoDb/EventSourcing/MongoEventStore_Reader.cs
  22. 7
      src/Squidex.Infrastructure/EventSourcing/StoredEvent.cs
  23. 12
      tests/Squidex.Domain.Apps.Entities.Tests/Backup/EventStreamTests.cs
  24. 12
      tests/Squidex.Infrastructure.Tests/EventSourcing/Grains/EventConsumerGrainTests.cs
  25. 4
      tests/Squidex.Infrastructure.Tests/EventSourcing/RetrySubscriptionTests.cs
  26. 4
      tests/Squidex.Infrastructure.Tests/States/PersistenceEventSourcingTests.cs

25
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++;

10
src/Squidex.Domain.Apps.Entities/Backup/EventStreamReader.cs

@ -34,7 +34,7 @@ namespace Squidex.Domain.Apps.Entities.Backup
}
}
public async Task ReadEventsAsync(Func<EventData, Stream, Task> eventHandler)
public async Task ReadEventsAsync(Func<StoredEvent, Stream, Task> 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++;
}
}
}

12
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<Stream, Task> attachment = null)
public async Task WriteEventAsync(StoredEvent storedEvent, Func<Stream, Task> 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++;
}
}
}

55
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<Guid> store;
protected HandlerBase(IStore<Guid> store)
{
Guard.NotNull(store, nameof(store));
this.store = store;
}
protected async Task RebuildManyAsync(IEnumerable<Guid> ids, Func<Guid, Task> action)
{
foreach (var id in ids)
{
await action(id);
}
}
protected async Task RebuildAsync<TState, TGrain>(Guid key, Func<Envelope<IEvent>, TState, TState> func) where TState : IDomainState, new()
{
var state = new TState
{
Version = EtagVersion.Empty
};
var persistence = store.WithSnapshotsAndEventSourcing<TState, Guid>(typeof(TGrain), key, s => state = s, e =>
{
state = func(e, state);
state.Version++;
});
await persistence.ReadAsync();
await persistence.WriteSnapshotAsync(state);
}
}
}

50
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<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

@ -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<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;
}
}
}

54
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<Guid> contentIds = new HashSet<Guid>();
public string Name { get; } = "Contents";
public RestoreContents(IStore<Guid> store)
: base(store)
{
}
public Task HandleAsync(Envelope<IEvent> @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<ContentState, ContentGrain>(id, (e, s) => s.Apply(e)));
}
public Task CompleteAsync()
{
return TaskHelper.Done;
}
}
}

70
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<Guid> ruleIds = new HashSet<Guid>();
private readonly IGrainFactory grainFactory;
private Guid appId;
public string Name { get; } = "Rules";
public RestoreRules(IStore<Guid> store, IGrainFactory grainFactory)
: base(store)
{
Guard.NotNull(grainFactory, nameof(grainFactory));
this.grainFactory = grainFactory;
}
public Task HandleAsync(Envelope<IEvent> @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<RuleState, RuleGrain>(id, (e, s) => s.Apply(e)));
await grainFactory.GetGrain<IRulesByAppIndex>(appId).RebuildAsync(ruleIds);
}
public Task CompleteAsync()
{
return TaskHelper.Done;
}
}
}

75
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<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)
: base(store)
{
Guard.NotNull(fieldRegistry, nameof(fieldRegistry));
Guard.NotNull(grainFactory, nameof(grainFactory));
this.fieldRegistry = fieldRegistry;
this.grainFactory = grainFactory;
}
public Task HandleAsync(Envelope<IEvent> @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<SchemaState, SchemaGrain>(id, (e, s) => s.Apply(e, fieldRegistry)));
await grainFactory.GetGrain<ISchemasByAppIndex>(appId).RebuildAsync(schemasByName);
}
public Task CompleteAsync()
{
return TaskHelper.Done;
}
}
}

21
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<J<IRestoreJob>> GetStateAsync();
}
}

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

23
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; }
}
}

203
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<string> store;
private readonly IEnumerable<IRestoreHandler> handlers;
private RestoreState state = new RestoreState();
private IPersistence<RestoreState> persistence;
public RestoreGrain(
IAssetStore assetStore,
IBackupArchiveLocation backupArchiveLocation,
IClock clock,
IEventStore eventStore,
IEventDataFormatter eventDataFormatter,
IGrainFactory grainFactory,
IEnumerable<IRestoreHandler> handlers,
ISemanticLog log,
IStore<Guid> 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<RestoreState, string>(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<EventData> { @event.Data });
eventIndex++;
state.Job.Status = $"Handled event {eventIndex}";
});
}
}
}
private async Task DoAsync(string start, string end, Func<Task> action)
{
state.Job.Status = start;
await action();
state.Job.Status = end;
}
public Task<J<IRestoreJob>> GetStateAsync()
{
return Task.FromResult<J<IRestoreJob>>(state.Job);
}
}
}

17
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; }
}
}

38
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; }
}
}

5
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<string, string> allTags)
{
return GetGrain(appId, group).RebuildTagsAsync(allTags);
}
public Task ClearAsync(Guid appId)
{
return GetGrain(appId, TagGroups.Assets).ClearAsync();

1
src/Squidex.Domain.Apps.Entities/Tags/ITagGrain.cs

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

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

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

5
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<string, string> allTags)
{
throw new NotImplementedException();
}
public async Task<HashSet<string>> NormalizeTagsAsync(HashSet<string> names, HashSet<string> ids)
{
var result = new HashSet<string>();

1
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);

4
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++;
}

7
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;
}
}
}

12
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)
{

12
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();

4
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<IEventSubscription>(), ev);
await sut.StopAsync();

4
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> { 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);

Loading…
Cancel
Save