From 0bfebccadbdd4b7482c946fb571162aac73e03dd Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sat, 3 Feb 2024 22:24:46 +0100 Subject: [PATCH] Add compression feature for workflow state data This update adds a compression feature for workflow state data to reduce storage needs. A compression strategy resolver, None and GZip strategies have been implemented. Migration scripts were also updated, and a superfluous file(s) were removed. --- Elsa.sln | 7 +- Elsa.sln.DotSettings | 1 + ....sh => generate-migrations-initial copy.sh | 2 +- migrations/efcore-3.0.sh | 38 +++++++++ migrations/efcore-3.1.sh | 35 ++++++++ .../Modules/Management/Configurations.cs | 2 + .../Management/WorkflowInstanceStore.cs | 80 ++++++++++++++----- .../Compression/GZip.cs | 36 +++++++++ .../Compression/None.cs | 21 +++++ .../Contracts/ICompressionStrategy.cs | 20 +++++ .../Contracts/ICompressionStrategyResolver.cs | 12 +++ .../Features/WorkflowManagementFeature.cs | 5 ++ .../Options/ManagementOptions.cs | 5 ++ .../Services/CompressionStrategyResolver.cs | 14 ++++ 14 files changed, 257 insertions(+), 21 deletions(-) rename update-migrations.sh => generate-migrations-initial copy.sh (98%) create mode 100644 migrations/efcore-3.0.sh create mode 100644 migrations/efcore-3.1.sh create mode 100644 src/modules/Elsa.Workflows.Management/Compression/GZip.cs create mode 100644 src/modules/Elsa.Workflows.Management/Compression/None.cs create mode 100644 src/modules/Elsa.Workflows.Management/Contracts/ICompressionStrategy.cs create mode 100644 src/modules/Elsa.Workflows.Management/Contracts/ICompressionStrategyResolver.cs create mode 100644 src/modules/Elsa.Workflows.Management/Services/CompressionStrategyResolver.cs diff --git a/Elsa.sln b/Elsa.sln index ffc51f4c2..3ae1a20f5 100644 --- a/Elsa.sln +++ b/Elsa.sln @@ -23,7 +23,6 @@ Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "solution", "solution", "{7D NuGet.Config = NuGet.Config packages.props = packages.props README.md = README.md - update-migrations.sh = update-migrations.sh EndProjectSection EndProject Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "docs", "docs", "{0354F050-3992-4DD4-B0EE-5FBA04AC72B6}" @@ -308,6 +307,12 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Samples.AspNet.Heartbe EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.MassTransit.AzureServiceBus", "src\modules\Elsa.MassTransit.AzureServiceBus\Elsa.MassTransit.AzureServiceBus.csproj", "{AFEB799E-82C3-4D02-9D5C-766BB8DEF004}" EndProject +Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "migrations", "migrations", "{C80C8231-D35C-4ACC-9ED6-9F3DB221535E}" + ProjectSection(SolutionItems) = preProject + migrations\efcore-3.1.sh = migrations\efcore-3.1.sh + migrations\efcore-3.0.sh = migrations\efcore-3.0.sh + EndProjectSection +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU diff --git a/Elsa.sln.DotSettings b/Elsa.sln.DotSettings index 05601a597..22b5924e3 100644 --- a/Elsa.sln.DotSettings +++ b/Elsa.sln.DotSettings @@ -3,6 +3,7 @@ 1 False True + EF True True True diff --git a/update-migrations.sh b/generate-migrations-initial copy.sh similarity index 98% rename from update-migrations.sh rename to generate-migrations-initial copy.sh index ee23f10a3..f5dacadae 100644 --- a/update-migrations.sh +++ b/generate-migrations-initial copy.sh @@ -1,7 +1,7 @@ #!/usr/bin/env zsh # Define the modules to update -mods=("Runtime") +mods=("Management") # mods=("Alterations" "Runtime" "Management" "Identity" "Labels") # Define the list of providers diff --git a/migrations/efcore-3.0.sh b/migrations/efcore-3.0.sh new file mode 100644 index 000000000..420d71497 --- /dev/null +++ b/migrations/efcore-3.0.sh @@ -0,0 +1,38 @@ +#!/usr/bin/env zsh + +# Define the modules to update +mods=("Management") +# mods=("Alterations" "Runtime" "Management" "Identity" "Labels") + +# Define the list of providers +providers=("MySql" "SqlServer" "Sqlite" "PostgreSql") +# providers=("SqlServer") + +# Connection strings for each provider +typeset -A connStrings +connStrings=( + MySql "Server=localhost;Port=3306;Database=elsa;User=root;Password=password;" + SqlServer "" + Sqlite "" + PostgreSql "" +) + +# Loop through each module +for module in "${mods[@]}"; do + # Loop through each provider + for provider in "${providers[@]}"; do + providerPath="../src/modules/Elsa.EntityFrameworkCore.$provider" + migrationsPath="Migrations/$module" + + echo "Updating migrations for $provider..." + echo "Provider path: ${providerPath:?}/${migrationsPath}" + echo "Migrations path: $migrationsPath" + echo "Connection string: ${connStrings[$provider]}" + + # 1. Delete the existing migrations folder + rm -rf "${providerPath:?}/${migrationsPath}" + + # 2. Run the migrations command + dotnet ef migrations add Initial -c "$module"ElsaDbContext -p "$providerPath" -o "$migrationsPath" -- --connectionString "${connStrings[$provider]}" + done +done diff --git a/migrations/efcore-3.1.sh b/migrations/efcore-3.1.sh new file mode 100644 index 000000000..82b3f4509 --- /dev/null +++ b/migrations/efcore-3.1.sh @@ -0,0 +1,35 @@ +#!/usr/bin/env zsh + +# Define the modules to update +mods=("Management") +# mods=("Alterations" "Runtime" "Management" "Identity" "Labels") + +# Define the list of providers +providers=("MySql" "SqlServer" "Sqlite" "PostgreSql") +# providers=("SqlServer") + +# Connection strings for each provider +typeset -A connStrings +connStrings=( + MySql "Server=localhost;Port=3306;Database=elsa;User=root;Password=password;" + SqlServer "" + Sqlite "" + PostgreSql "" +) + +# Loop through each module +for module in "${mods[@]}"; do + # Loop through each provider + for provider in "${providers[@]}"; do + providerPath="../src/modules/Elsa.EntityFrameworkCore.$provider" + migrationsPath="Migrations/$module" + + echo "Updating migrations for $provider..." + echo "Provider path: ${providerPath:?}/${migrationsPath}" + echo "Migrations path: $migrationsPath" + echo "Connection string: ${connStrings[$provider]}" + + # 1. Run the migrations command + dotnet ef migrations add 3_1 -c "$module"ElsaDbContext -p "$providerPath" -o "$migrationsPath" -- --connectionString "${connStrings[$provider]}" + done +done diff --git a/src/modules/Elsa.EntityFrameworkCore/Modules/Management/Configurations.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Management/Configurations.cs index f2193ea2b..547ba59df 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Management/Configurations.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Management/Configurations.cs @@ -19,6 +19,8 @@ internal class Configurations : IEntityTypeConfiguration, IE builder.Ignore(x => x.CustomProperties); builder.Ignore(x => x.Options); builder.Property("Data"); + builder.Property("DataFormat"); + builder.Property("DataCompressionAlgorithm"); builder.Property("UsableAsActivity"); builder.Property(x => x.ToolVersion).HasConversion(VersionToStringConverter, StringToVersionConverter); diff --git a/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs index 01f8e8dbd..814257271 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs @@ -2,11 +2,15 @@ using Elsa.Common.Models; using Elsa.EntityFrameworkCore.Common; using Elsa.Extensions; using Elsa.Workflows.Contracts; +using Elsa.Workflows.Management.Compression; using Elsa.Workflows.Management.Contracts; using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Management.Filters; using Elsa.Workflows.Management.Models; +using Elsa.Workflows.Management.Options; +using JetBrains.Annotations; using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.Options; using Open.Linq.AsyncExtensions; namespace Elsa.EntityFrameworkCore.Modules.Management; @@ -14,23 +18,34 @@ namespace Elsa.EntityFrameworkCore.Modules.Management; /// /// An EF Core implementation of . /// +[UsedImplicitly] public class EFCoreWorkflowInstanceStore : IWorkflowInstanceStore { private readonly EntityStore _store; private readonly IWorkflowStateSerializer _workflowStateSerializer; + private readonly ICompressionStrategyResolver _compressionStrategyResolver; + private readonly IOptions _options; /// /// Constructor. /// - public EFCoreWorkflowInstanceStore(EntityStore store, IWorkflowStateSerializer workflowStateSerializer) + public EFCoreWorkflowInstanceStore( + EntityStore store, + IWorkflowStateSerializer workflowStateSerializer, + ICompressionStrategyResolver compressionStrategyResolver, + IOptions options) { _store = store; _workflowStateSerializer = workflowStateSerializer; + _compressionStrategyResolver = compressionStrategyResolver; + _options = options; } /// - public async ValueTask FindAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default) => - await _store.QueryAsync(query => Filter(query, filter), OnLoadAsync, cancellationToken).FirstOrDefault(); + public async ValueTask FindAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default) + { + return await _store.QueryAsync(query => Filter(query, filter), OnLoadAsync, cancellationToken).FirstOrDefault(); + } /// public async ValueTask> FindManyAsync(WorkflowInstanceFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default) @@ -49,13 +64,17 @@ public class EFCoreWorkflowInstanceStore : IWorkflowInstanceStore } /// - public async ValueTask> FindManyAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default) => - await _store.QueryAsync(query => Filter(query, filter), OnLoadAsync, cancellationToken).ToList().AsEnumerable(); + public async ValueTask> FindManyAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default) + { + return await _store.QueryAsync(query => Filter(query, filter), OnLoadAsync, cancellationToken).ToList().AsEnumerable(); + } /// - public async ValueTask> FindManyAsync(WorkflowInstanceFilter filter, WorkflowInstanceOrder order, CancellationToken cancellationToken = default) => - await _store.QueryAsync(query => Filter(query, filter).OrderBy(order), OnLoadAsync, cancellationToken).ToList().AsEnumerable(); - + public async ValueTask> FindManyAsync(WorkflowInstanceFilter filter, WorkflowInstanceOrder order, CancellationToken cancellationToken = default) + { + return await _store.QueryAsync(query => Filter(query, filter).OrderBy(order), OnLoadAsync, cancellationToken).ToList().AsEnumerable(); + } + /// public async ValueTask CountAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default) { @@ -84,31 +103,46 @@ public class EFCoreWorkflowInstanceStore : IWorkflowInstanceStore } /// - public async ValueTask> SummarizeManyAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default) => - await _store.QueryAsync(query => Filter(query, filter), WorkflowInstanceSummary.FromInstanceExpression(), cancellationToken).ToList().AsEnumerable(); + public async ValueTask> SummarizeManyAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default) + { + return await _store.QueryAsync(query => Filter(query, filter), WorkflowInstanceSummary.FromInstanceExpression(), cancellationToken).ToList().AsEnumerable(); + } /// - public async ValueTask> SummarizeManyAsync(WorkflowInstanceFilter filter, WorkflowInstanceOrder order, CancellationToken cancellationToken = default) => - await _store.QueryAsync(query => Filter(query, filter).OrderBy(order), WorkflowInstanceSummary.FromInstanceExpression(), cancellationToken).ToList().AsEnumerable(); + public async ValueTask> SummarizeManyAsync(WorkflowInstanceFilter filter, WorkflowInstanceOrder order, CancellationToken cancellationToken = default) + { + return await _store.QueryAsync(query => Filter(query, filter).OrderBy(order), WorkflowInstanceSummary.FromInstanceExpression(), cancellationToken).ToList().AsEnumerable(); + } /// - public async ValueTask DeleteAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default) => - await _store.DeleteWhereAsync(query => Filter(query, filter), cancellationToken); + public async ValueTask DeleteAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default) + { + return await _store.DeleteWhereAsync(query => Filter(query, filter), cancellationToken); + } /// - public async ValueTask SaveAsync(WorkflowInstance instance, CancellationToken cancellationToken = default) => + public async ValueTask SaveAsync(WorkflowInstance instance, CancellationToken cancellationToken = default) + { await _store.SaveAsync(instance, OnSaveAsync, cancellationToken); + } /// - public async ValueTask SaveManyAsync(IEnumerable instances, CancellationToken cancellationToken = default) => + public async ValueTask SaveManyAsync(IEnumerable instances, CancellationToken cancellationToken = default) + { await _store.SaveManyAsync(instances, OnSaveAsync, cancellationToken); + } private async ValueTask OnSaveAsync(ManagementElsaDbContext managementElsaDbContext, WorkflowInstance entity, CancellationToken cancellationToken) { var data = entity.WorkflowState; var json = await _workflowStateSerializer.SerializeAsync(data, cancellationToken); + var compressionAlgorithm = _options.Value.CompressionAlgorithm ?? nameof(None); + var compressionStrategy = _compressionStrategyResolver.Resolve(compressionAlgorithm); + var compressedJson = await compressionStrategy.CompressAsync(json, cancellationToken); - managementElsaDbContext.Entry(entity).Property("Data").CurrentValue = json; + managementElsaDbContext.Entry(entity).Property("Data").CurrentValue = compressedJson; + managementElsaDbContext.Entry(entity).Property("DataFormat").CurrentValue = "Json"; + managementElsaDbContext.Entry(entity).Property("DataCompressionAlgorithm").CurrentValue = compressionAlgorithm; } private async ValueTask OnLoadAsync(ManagementElsaDbContext managementElsaDbContext, WorkflowInstance? entity, CancellationToken cancellationToken) @@ -118,12 +152,20 @@ public class EFCoreWorkflowInstanceStore : IWorkflowInstanceStore var data = entity.WorkflowState; var json = (string?)managementElsaDbContext.Entry(entity).Property("Data").CurrentValue; + var compressionAlgorithm = (string?)managementElsaDbContext.Entry(entity).Property("DataCompressionAlgorithm").CurrentValue ?? nameof(None); + var compressionStrategy = _compressionStrategyResolver.Resolve(compressionAlgorithm); - if (!string.IsNullOrWhiteSpace(json)) + if (!string.IsNullOrWhiteSpace(json)) + { + json = await compressionStrategy.DecompressAsync(json, cancellationToken); data = await _workflowStateSerializer.DeserializeAsync(json, cancellationToken); + } entity.WorkflowState = data; } - private static IQueryable Filter(IQueryable query, WorkflowInstanceFilter filter) => filter.Apply(query); + private static IQueryable Filter(IQueryable query, WorkflowInstanceFilter filter) + { + return filter.Apply(query); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Compression/GZip.cs b/src/modules/Elsa.Workflows.Management/Compression/GZip.cs new file mode 100644 index 000000000..5dba3da21 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Compression/GZip.cs @@ -0,0 +1,36 @@ +using System.IO.Compression; +using System.Text; +using Elsa.Workflows.Management.Contracts; + +namespace Elsa.Workflows.Management.Compression; + +/// +/// Represents a GZip compression strategy. +/// +public class GZip : ICompressionStrategy +{ + /// + public async ValueTask CompressAsync(string input, CancellationToken cancellationToken) + { + var inputBytes = Encoding.UTF8.GetBytes(input); + using var output = new MemoryStream(); + await using var compressionStream = new GZipStream(output, CompressionMode.Compress); + await compressionStream.WriteAsync(inputBytes, 0, inputBytes.Length, cancellationToken); + + return Convert.ToBase64String(output.ToArray()); + } + + /// + public async ValueTask DecompressAsync(string input, CancellationToken cancellationToken) + { + var inputBytes = Convert.FromBase64String(input); + using var inputMemoryStream = new MemoryStream(inputBytes); + await using var decompressionStream = new GZipStream(inputMemoryStream, CompressionMode.Decompress); + using var outputMemoryStream = new MemoryStream(); + + await decompressionStream.CopyToAsync(outputMemoryStream, cancellationToken); + var decompressedBytes = outputMemoryStream.ToArray(); + + return Encoding.UTF8.GetString(decompressedBytes); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Compression/None.cs b/src/modules/Elsa.Workflows.Management/Compression/None.cs new file mode 100644 index 000000000..3072c74b3 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Compression/None.cs @@ -0,0 +1,21 @@ +using Elsa.Workflows.Management.Contracts; + +namespace Elsa.Workflows.Management.Compression; + +/// +/// Represents a compression strategy that does not compress or decompress the input. +/// +public class None : ICompressionStrategy +{ + /// + public ValueTask CompressAsync(string input, CancellationToken cancellationToken) + { + return new ValueTask(input); + } + + /// + public ValueTask DecompressAsync(string input, CancellationToken cancellationToken) + { + return new ValueTask(input); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Contracts/ICompressionStrategy.cs b/src/modules/Elsa.Workflows.Management/Contracts/ICompressionStrategy.cs new file mode 100644 index 000000000..150611459 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Contracts/ICompressionStrategy.cs @@ -0,0 +1,20 @@ +namespace Elsa.Workflows.Management.Contracts; + +/// +/// Represents a compression strategy. +/// +public interface ICompressionStrategy +{ + /// + /// Compresses the input. + /// + ValueTask CompressAsync(string input, CancellationToken cancellationToken); + + /// + /// Decompresses the input. + /// + /// + /// + /// + ValueTask DecompressAsync(string input, CancellationToken cancellationToken); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Contracts/ICompressionStrategyResolver.cs b/src/modules/Elsa.Workflows.Management/Contracts/ICompressionStrategyResolver.cs new file mode 100644 index 000000000..df7824b3c --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Contracts/ICompressionStrategyResolver.cs @@ -0,0 +1,12 @@ +namespace Elsa.Workflows.Management.Contracts; + +/// +/// Resolves a from its name. +/// +public interface ICompressionStrategyResolver +{ + /// + /// Resolves a from its name. + /// + ICompressionStrategy Resolve(string name); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs index fb428d9b1..18971ff3b 100644 --- a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs +++ b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs @@ -1,5 +1,6 @@ using System.ComponentModel; using System.Dynamic; +using System.IO.Compression; using System.Reflection; using Elsa.Common.Contracts; using Elsa.Common.Features; @@ -11,6 +12,7 @@ using Elsa.Features.Services; using Elsa.Workflows.Contracts; using Elsa.Workflows.Features; using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; +using Elsa.Workflows.Management.Compression; using Elsa.Workflows.Management.Contracts; using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Management.Mappers; @@ -175,6 +177,9 @@ public class WorkflowManagementFeature : FeatureBase .AddScoped() .AddSingleton() .AddSingleton() + .AddSingleton() + .AddSingleton() + .AddSingleton() ; Services.AddNotificationHandlersFrom(GetType()); diff --git a/src/modules/Elsa.Workflows.Management/Options/ManagementOptions.cs b/src/modules/Elsa.Workflows.Management/Options/ManagementOptions.cs index 0a65f8498..200c5557d 100644 --- a/src/modules/Elsa.Workflows.Management/Options/ManagementOptions.cs +++ b/src/modules/Elsa.Workflows.Management/Options/ManagementOptions.cs @@ -16,4 +16,9 @@ public class ManagementOptions /// A collection of types that are available to the system as variable types. /// public HashSet VariableDescriptors { get; set; } = new(); + + /// + /// The format to use for compressing workflow state. + /// + public string? CompressionAlgorithm { get; set; } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Services/CompressionStrategyResolver.cs b/src/modules/Elsa.Workflows.Management/Services/CompressionStrategyResolver.cs new file mode 100644 index 000000000..bb435aa75 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Services/CompressionStrategyResolver.cs @@ -0,0 +1,14 @@ +using Elsa.Workflows.Management.Compression; +using Elsa.Workflows.Management.Contracts; + +namespace Elsa.Workflows.Management.Services; + +/// +public class CompressionStrategyResolver(IEnumerable strategies) : ICompressionStrategyResolver +{ + /// + public ICompressionStrategy Resolve(string name) + { + return strategies.FirstOrDefault(s => s.GetType().Name == name) ?? new None(); + } +} \ No newline at end of file