mirror of https://github.com/Squidex/squidex.git
10 changed files with 4 additions and 459 deletions
@ -1,242 +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.Linq; |
|
||||
using System.Threading; |
|
||||
using System.Threading.Tasks; |
|
||||
using NodaTime; |
|
||||
using Orleans; |
|
||||
using Orleans.Concurrency; |
|
||||
using Squidex.Domain.Apps.Backup.State; |
|
||||
using Squidex.Domain.Apps.Events.Assets; |
|
||||
using Squidex.Infrastructure; |
|
||||
using Squidex.Infrastructure.Assets; |
|
||||
using Squidex.Infrastructure.EventSourcing; |
|
||||
using Squidex.Infrastructure.Log; |
|
||||
using Squidex.Infrastructure.Orleans; |
|
||||
using Squidex.Infrastructure.States; |
|
||||
|
|
||||
namespace Squidex.Domain.Apps.Backup |
|
||||
{ |
|
||||
[Reentrant] |
|
||||
public sealed class BackupGrain : Grain, IBackupGrain |
|
||||
{ |
|
||||
private const int MaxBackups = 10; |
|
||||
private readonly IClock clock; |
|
||||
private readonly IAssetStore assetStore; |
|
||||
private readonly IEventDataFormatter eventDataFormatter; |
|
||||
private readonly ISemanticLog log; |
|
||||
private readonly IEventStore eventStore; |
|
||||
private readonly IBackupArchiveLocation backupArchiveLocation; |
|
||||
private readonly IStore<Guid> store; |
|
||||
private CancellationTokenSource currentTask; |
|
||||
private BackupStateJob currentJob; |
|
||||
private Guid appId; |
|
||||
private BackupState state; |
|
||||
private IPersistence<BackupState> persistence; |
|
||||
|
|
||||
public BackupGrain( |
|
||||
IAssetStore assetStore, |
|
||||
IBackupArchiveLocation backupArchiveLocation, |
|
||||
IClock clock, |
|
||||
IEventStore eventStore, |
|
||||
IEventDataFormatter eventDataFormatter, |
|
||||
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(store, nameof(store)); |
|
||||
Guard.NotNull(log, nameof(log)); |
|
||||
|
|
||||
this.assetStore = assetStore; |
|
||||
this.backupArchiveLocation = backupArchiveLocation; |
|
||||
this.clock = clock; |
|
||||
this.eventStore = eventStore; |
|
||||
this.eventDataFormatter = eventDataFormatter; |
|
||||
this.store = store; |
|
||||
this.log = log; |
|
||||
} |
|
||||
|
|
||||
public override Task OnActivateAsync() |
|
||||
{ |
|
||||
return OnActivateAsync(this.GetPrimaryKey()); |
|
||||
} |
|
||||
|
|
||||
public async Task OnActivateAsync(Guid appId) |
|
||||
{ |
|
||||
this.appId = appId; |
|
||||
|
|
||||
persistence = store.WithSnapshots<BackupState, Guid>(GetType(), appId, s => state = s); |
|
||||
|
|
||||
await ReadAsync(); |
|
||||
await CleanupAsync(); |
|
||||
} |
|
||||
|
|
||||
private async Task ReadAsync() |
|
||||
{ |
|
||||
await persistence.ReadAsync(); |
|
||||
} |
|
||||
|
|
||||
private async Task WriteAsync() |
|
||||
{ |
|
||||
await persistence.WriteSnapshotAsync(state); |
|
||||
} |
|
||||
|
|
||||
private async Task CleanupAsync() |
|
||||
{ |
|
||||
var hasUpdated = false; |
|
||||
|
|
||||
foreach (var job in state.Jobs) |
|
||||
{ |
|
||||
if (!job.Stopped.HasValue) |
|
||||
{ |
|
||||
await CleanupAsync(job); |
|
||||
|
|
||||
job.Stopped = clock.GetCurrentInstant(); |
|
||||
job.Failed = true; |
|
||||
|
|
||||
hasUpdated = true; |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
if (hasUpdated) |
|
||||
{ |
|
||||
await WriteAsync(); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private async Task CleanupAsync(BackupStateJob job) |
|
||||
{ |
|
||||
await backupArchiveLocation.DeleteArchiveAsync(job.Id); |
|
||||
} |
|
||||
|
|
||||
public async Task StartNewAsync() |
|
||||
{ |
|
||||
if (currentTask != null) |
|
||||
{ |
|
||||
throw new DomainException("Another backup process is already running."); |
|
||||
} |
|
||||
|
|
||||
if (state.Jobs.Count >= MaxBackups) |
|
||||
{ |
|
||||
throw new DomainException($"You cannot have more than {MaxBackups} backups."); |
|
||||
} |
|
||||
|
|
||||
var job = new BackupStateJob { Id = Guid.NewGuid(), Started = clock.GetCurrentInstant() }; |
|
||||
|
|
||||
currentTask = new CancellationTokenSource(); |
|
||||
currentJob = job; |
|
||||
|
|
||||
state.Jobs.Add(job); |
|
||||
|
|
||||
await WriteAsync(); |
|
||||
|
|
||||
try |
|
||||
{ |
|
||||
using (var stream = await backupArchiveLocation.OpenStreamAsync(job.Id)) |
|
||||
{ |
|
||||
using (var writer = new EventStreamWriter(stream)) |
|
||||
{ |
|
||||
await eventStore.QueryAsync(async @event => |
|
||||
{ |
|
||||
var eventData = @event.Data; |
|
||||
|
|
||||
if (eventData.Type == nameof(AssetCreated) || |
|
||||
eventData.Type == nameof(AssetUpdated)) |
|
||||
{ |
|
||||
var parsedEvent = eventDataFormatter.Parse(eventData); |
|
||||
|
|
||||
var assetVersion = 0L; |
|
||||
var assetId = Guid.Empty; |
|
||||
|
|
||||
if (parsedEvent.Payload is AssetCreated assetCreated) |
|
||||
{ |
|
||||
assetId = assetCreated.AssetId; |
|
||||
assetVersion = assetCreated.FileVersion; |
|
||||
} |
|
||||
|
|
||||
if (parsedEvent.Payload is AssetUpdated asetUpdated) |
|
||||
{ |
|
||||
assetId = asetUpdated.AssetId; |
|
||||
assetVersion = asetUpdated.FileVersion; |
|
||||
} |
|
||||
|
|
||||
await writer.WriteEventAsync(eventData, async attachmentStream => |
|
||||
{ |
|
||||
await assetStore.DownloadAsync(assetId.ToString(), assetVersion, null, attachmentStream); |
|
||||
}); |
|
||||
} |
|
||||
else |
|
||||
{ |
|
||||
await writer.WriteEventAsync(eventData); |
|
||||
} |
|
||||
}, "AppId", appId, null, currentTask.Token); |
|
||||
} |
|
||||
|
|
||||
stream.Position = 0; |
|
||||
|
|
||||
currentTask.Token.ThrowIfCancellationRequested(); |
|
||||
|
|
||||
await assetStore.UploadAsync(job.Id.ToString(), 0, null, stream); |
|
||||
|
|
||||
currentTask.Token.ThrowIfCancellationRequested(); |
|
||||
} |
|
||||
} |
|
||||
catch |
|
||||
{ |
|
||||
job.Failed = true; |
|
||||
} |
|
||||
finally |
|
||||
{ |
|
||||
job.Stopped = clock.GetCurrentInstant(); |
|
||||
|
|
||||
await WriteAsync(); |
|
||||
|
|
||||
currentTask = null; |
|
||||
currentJob = null; |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
public async Task DeleteAsync(Guid id) |
|
||||
{ |
|
||||
var job = state.Jobs.FirstOrDefault(x => x.Id == id); |
|
||||
|
|
||||
if (job == null) |
|
||||
{ |
|
||||
throw new DomainObjectNotFoundException(id.ToString(), typeof(IBackupJob)); |
|
||||
} |
|
||||
|
|
||||
if (currentJob == job) |
|
||||
{ |
|
||||
currentTask?.Cancel(); |
|
||||
} |
|
||||
else |
|
||||
{ |
|
||||
state.Jobs.Remove(job); |
|
||||
|
|
||||
await WriteAsync(); |
|
||||
await CleanupAsync(job); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
public Task<J<List<IBackupJob>>> GetStateAsync() |
|
||||
{ |
|
||||
return Task.FromResult(new J<List<IBackupJob>>(state.Jobs.OfType<IBackupJob>().ToList())); |
|
||||
} |
|
||||
|
|
||||
private bool IsRunning() |
|
||||
{ |
|
||||
return state.Jobs.Any(x => !x.Stopped.HasValue); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
@ -1,80 +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.IO.Compression; |
|
||||
using System.Threading.Tasks; |
|
||||
using Newtonsoft.Json; |
|
||||
using Newtonsoft.Json.Linq; |
|
||||
using Squidex.Infrastructure; |
|
||||
using Squidex.Infrastructure.EventSourcing; |
|
||||
|
|
||||
namespace Squidex.Domain.Apps.Backup |
|
||||
{ |
|
||||
public sealed class EventStreamWriter : DisposableObjectBase |
|
||||
{ |
|
||||
private const int MaxItemsPerFolder = 1000; |
|
||||
private readonly StreamWriter streamWriter; |
|
||||
private readonly ZipArchive archive; |
|
||||
private int writtenEvents; |
|
||||
private int writtenAttachments; |
|
||||
|
|
||||
public EventStreamWriter(Stream stream) |
|
||||
{ |
|
||||
archive = new ZipArchive(stream, ZipArchiveMode.Update, true); |
|
||||
} |
|
||||
|
|
||||
public async Task WriteEventAsync(EventData eventData, 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 eventFolder = writtenEvents / MaxItemsPerFolder; |
|
||||
var eventPath = $"events/{eventFolder}/{writtenEvents}.json"; |
|
||||
var eventEntry = archive.GetEntry(eventPath) ?? archive.CreateEntry(eventPath); |
|
||||
|
|
||||
using (var stream = eventEntry.Open()) |
|
||||
{ |
|
||||
using (var textWriter = new StreamWriter(stream)) |
|
||||
{ |
|
||||
using (var jsonWriter = new JsonTextWriter(textWriter)) |
|
||||
{ |
|
||||
await eventObject.WriteToAsync(jsonWriter); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
writtenEvents++; |
|
||||
|
|
||||
if (attachment != null) |
|
||||
{ |
|
||||
var attachmentFolder = writtenAttachments / MaxItemsPerFolder; |
|
||||
var attachmentPath = $"attachments/{attachmentFolder}/{writtenEvents}.blob"; |
|
||||
var attachmentEntry = archive.GetEntry(attachmentPath) ?? archive.CreateEntry(attachmentPath); |
|
||||
|
|
||||
using (var stream = eventEntry.Open()) |
|
||||
{ |
|
||||
await attachment(stream); |
|
||||
} |
|
||||
|
|
||||
writtenAttachments++; |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
protected override void DisposeObject(bool disposing) |
|
||||
{ |
|
||||
if (disposing) |
|
||||
{ |
|
||||
archive.Dispose(); |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
} |
|
||||
@ -1,13 +0,0 @@ |
|||||
using System; |
|
||||
using System.IO; |
|
||||
using System.Threading.Tasks; |
|
||||
|
|
||||
namespace Squidex.Domain.Apps.Backup |
|
||||
{ |
|
||||
public interface IBackupArchiveLocation |
|
||||
{ |
|
||||
Task<Stream> OpenStreamAsync(Guid backupId); |
|
||||
|
|
||||
Task DeleteArchiveAsync(Guid backupId); |
|
||||
} |
|
||||
} |
|
||||
@ -1,24 +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.Threading.Tasks; |
|
||||
using Orleans; |
|
||||
using Squidex.Infrastructure.Orleans; |
|
||||
|
|
||||
namespace Squidex.Domain.Apps.Backup |
|
||||
{ |
|
||||
public interface IBackupGrain : IGrainWithGuidKey |
|
||||
{ |
|
||||
Task StartNewAsync(); |
|
||||
|
|
||||
Task DeleteAsync(Guid id); |
|
||||
|
|
||||
Task<J<List<IBackupJob>>> GetStateAsync(); |
|
||||
} |
|
||||
} |
|
||||
@ -1,25 +0,0 @@ |
|||||
// ==========================================================================
|
|
||||
// Squidex Headless CMS
|
|
||||
// ==========================================================================
|
|
||||
// Copyright (c) Squidex UG (haftungsbeschraenkt)
|
|
||||
// All rights reserved. Licensed under the MIT license.
|
|
||||
// ==========================================================================
|
|
||||
|
|
||||
using System; |
|
||||
using NodaTime; |
|
||||
|
|
||||
namespace Squidex.Domain.Apps.Backup |
|
||||
{ |
|
||||
public interface IBackupJob |
|
||||
{ |
|
||||
Guid Id { get; } |
|
||||
|
|
||||
Instant Started { get; } |
|
||||
|
|
||||
Instant? Stopped { get; } |
|
||||
|
|
||||
bool Failed { get; } |
|
||||
|
|
||||
string DownloadPath { get; } |
|
||||
} |
|
||||
} |
|
||||
@ -1,22 +0,0 @@ |
|||||
<Project Sdk="Microsoft.NET.Sdk"> |
|
||||
<PropertyGroup> |
|
||||
<TargetFramework>netstandard2.0</TargetFramework> |
|
||||
</PropertyGroup> |
|
||||
<ItemGroup> |
|
||||
<ProjectReference Include="..\Squidex.Domain.Apps.Entities\Squidex.Domain.Apps.Entities.csproj" /> |
|
||||
<ProjectReference Include="..\Squidex.Domain.Apps.Events\Squidex.Domain.Apps.Events.csproj" /> |
|
||||
</ItemGroup> |
|
||||
<ItemGroup> |
|
||||
<PackageReference Include="Microsoft.Orleans.OrleansCodeGenerator.Build" Version="2.0.0-rc2" /> |
|
||||
<PackageReference Include="NodaTime" Version="2.2.4" /> |
|
||||
<PackageReference Include="RefactoringEssentials" Version="5.6.0" /> |
|
||||
<PackageReference Include="StyleCop.Analyzers" Version="1.0.2" /> |
|
||||
<PackageReference Include="System.ValueTuple" Version="4.4.0" /> |
|
||||
</ItemGroup> |
|
||||
<PropertyGroup> |
|
||||
<CodeAnalysisRuleSet>..\..\Squidex.ruleset</CodeAnalysisRuleSet> |
|
||||
</PropertyGroup> |
|
||||
<ItemGroup> |
|
||||
<AdditionalFiles Include="..\..\stylecop.json" Link="stylecop.json" /> |
|
||||
</ItemGroup> |
|
||||
</Project> |
|
||||
@ -1,18 +0,0 @@ |
|||||
// ==========================================================================
|
|
||||
// Squidex Headless CMS
|
|
||||
// ==========================================================================
|
|
||||
// Copyright (c) Squidex UG (haftungsbeschraenkt)
|
|
||||
// All rights reserved. Licensed under the MIT license.
|
|
||||
// ==========================================================================
|
|
||||
|
|
||||
using System.Collections.Generic; |
|
||||
using Newtonsoft.Json; |
|
||||
|
|
||||
namespace Squidex.Domain.Apps.Backup.State |
|
||||
{ |
|
||||
public sealed class BackupState |
|
||||
{ |
|
||||
[JsonProperty] |
|
||||
public List<BackupStateJob> Jobs { get; set; } = new List<BackupStateJob>(); |
|
||||
} |
|
||||
} |
|
||||
@ -1,31 +0,0 @@ |
|||||
// ==========================================================================
|
|
||||
// Squidex Headless CMS
|
|
||||
// ==========================================================================
|
|
||||
// Copyright (c) Squidex UG (haftungsbeschraenkt)
|
|
||||
// All rights reserved. Licensed under the MIT license.
|
|
||||
// ==========================================================================
|
|
||||
|
|
||||
using System; |
|
||||
using Newtonsoft.Json; |
|
||||
using NodaTime; |
|
||||
|
|
||||
namespace Squidex.Domain.Apps.Backup.State |
|
||||
{ |
|
||||
public sealed class BackupStateJob : IBackupJob |
|
||||
{ |
|
||||
[JsonProperty] |
|
||||
public Guid Id { get; set; } |
|
||||
|
|
||||
[JsonProperty] |
|
||||
public Instant Started { get; set; } |
|
||||
|
|
||||
[JsonProperty] |
|
||||
public Instant? Stopped { get; set; } |
|
||||
|
|
||||
[JsonProperty] |
|
||||
public string DownloadPath { get; set; } |
|
||||
|
|
||||
[JsonProperty] |
|
||||
public bool Failed { get; set; } |
|
||||
} |
|
||||
} |
|
||||
Loading…
Reference in new issue