From 611c72fa266ca859e2952dff90fdd45e51f5afca Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sun, 4 Feb 2024 12:18:55 +0100 Subject: [PATCH] Implement Zstd compression and update related components Removed ICompressionStrategyResolver interface and file, and added ICompressionCodec and ICompressionCodecResolver interface. Updated relevant classes for the new interface. Specifically, added Zstd class under Compression as a new compression method. Also modified WorkflowInstanceStore.cs, None.cs, EFCoreWorkflowInstanceStore.cs, GZip.cs, and ActivityExecutionLogStore.cs for uniform compression terminology. --- Elsa.sln.DotSettings | 5 ++- src/bundles/Elsa.Server.Web/Program.cs | 2 +- .../Management/WorkflowInstanceStore.cs | 12 +++--- .../Runtime/ActivityExecutionLogStore.cs | 10 ++--- .../Compression/GZip.cs | 2 +- .../Compression/None.cs | 2 +- .../Compression/Zstd.cs | 37 +++++++++++++++++++ ...essionStrategy.cs => ICompressionCodec.cs} | 2 +- .../Contracts/ICompressionCodecResolver.cs | 12 ++++++ .../Contracts/ICompressionStrategyResolver.cs | 12 ------ .../Elsa.Workflows.Management.csproj | 1 + .../Features/WorkflowManagementFeature.cs | 7 ++-- .../Services/CompressionCodecResolver.cs | 24 ++++++++++++ .../Services/CompressionStrategyResolver.cs | 14 ------- 14 files changed, 96 insertions(+), 46 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Management/Compression/Zstd.cs rename src/modules/Elsa.Workflows.Management/Contracts/{ICompressionStrategy.cs => ICompressionCodec.cs} (92%) create mode 100644 src/modules/Elsa.Workflows.Management/Contracts/ICompressionCodecResolver.cs delete mode 100644 src/modules/Elsa.Workflows.Management/Contracts/ICompressionStrategyResolver.cs create mode 100644 src/modules/Elsa.Workflows.Management/Services/CompressionCodecResolver.cs delete mode 100644 src/modules/Elsa.Workflows.Management/Services/CompressionStrategyResolver.cs diff --git a/Elsa.sln.DotSettings b/Elsa.sln.DotSettings index 8bb0ccf33..d88a97a28 100644 --- a/Elsa.sln.DotSettings +++ b/Elsa.sln.DotSettings @@ -1,4 +1,4 @@ - + False 1000 1 @@ -17,4 +17,5 @@ True True True - True \ No newline at end of file + True + True \ No newline at end of file diff --git a/src/bundles/Elsa.Server.Web/Program.cs b/src/bundles/Elsa.Server.Web/Program.cs index aaef7368d..64b000aa7 100644 --- a/src/bundles/Elsa.Server.Web/Program.cs +++ b/src/bundles/Elsa.Server.Web/Program.cs @@ -110,7 +110,7 @@ services }); if(useZipCompression) - management.SetCompressionAlgorithm(nameof(GZip)); + management.SetCompressionAlgorithm(nameof(Zstd)); }) .UseWorkflowRuntime(runtime => { diff --git a/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs index b1df465b2..6990b94d6 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Management/WorkflowInstanceStore.cs @@ -23,7 +23,7 @@ public class EFCoreWorkflowInstanceStore : IWorkflowInstanceStore { private readonly EntityStore _store; private readonly IWorkflowStateSerializer _workflowStateSerializer; - private readonly ICompressionStrategyResolver _compressionStrategyResolver; + private readonly ICompressionCodecResolver _compressionCodecResolver; private readonly IOptions _options; /// @@ -32,12 +32,12 @@ public class EFCoreWorkflowInstanceStore : IWorkflowInstanceStore public EFCoreWorkflowInstanceStore( EntityStore store, IWorkflowStateSerializer workflowStateSerializer, - ICompressionStrategyResolver compressionStrategyResolver, + ICompressionCodecResolver compressionCodecResolver, IOptions options) { _store = store; _workflowStateSerializer = workflowStateSerializer; - _compressionStrategyResolver = compressionStrategyResolver; + _compressionCodecResolver = compressionCodecResolver; _options = options; } @@ -137,8 +137,8 @@ public class EFCoreWorkflowInstanceStore : IWorkflowInstanceStore 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); + var compressionCodec = _compressionCodecResolver.Resolve(compressionAlgorithm); + var compressedJson = await compressionCodec.CompressAsync(json, cancellationToken); managementElsaDbContext.Entry(entity).Property("Data").CurrentValue = compressedJson; managementElsaDbContext.Entry(entity).Property("DataCompressionAlgorithm").CurrentValue = compressionAlgorithm; @@ -152,7 +152,7 @@ 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); + var compressionStrategy = _compressionCodecResolver.Resolve(compressionAlgorithm); if (!string.IsNullOrWhiteSpace(json)) { diff --git a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs index d0402efb7..4eb12f62a 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs @@ -25,7 +25,7 @@ public class EFCoreActivityExecutionStore : IActivityExecutionStore private readonly EntityStore _store; private readonly ISafeSerializer _safeSerializer; private readonly IPayloadSerializer _payloadSerializer; - private readonly ICompressionStrategyResolver _compressionStrategyResolver; + private readonly ICompressionCodecResolver _compressionCodecResolver; private readonly IOptions _options; /// @@ -35,13 +35,13 @@ public class EFCoreActivityExecutionStore : IActivityExecutionStore EntityStore store, ISafeSerializer safeSerializer, IPayloadSerializer payloadSerializer, - ICompressionStrategyResolver compressionStrategyResolver, + ICompressionCodecResolver compressionCodecResolver, IOptions options) { _store = store; _safeSerializer = safeSerializer; _payloadSerializer = payloadSerializer; - _compressionStrategyResolver = compressionStrategyResolver; + _compressionCodecResolver = compressionCodecResolver; _options = options; } @@ -75,7 +75,7 @@ public class EFCoreActivityExecutionStore : IActivityExecutionStore { var compressionAlgorithm = _options.Value.CompressionAlgorithm ?? nameof(None); var serializedActivityState = entity.ActivityState != null ? await _safeSerializer.SerializeAsync(entity.ActivityState, cancellationToken) : default; - var compressedSerializedActivityState = serializedActivityState != null ? await _compressionStrategyResolver.Resolve(compressionAlgorithm).CompressAsync(serializedActivityState, cancellationToken) : default; + var compressedSerializedActivityState = serializedActivityState != null ? await _compressionCodecResolver.Resolve(compressionAlgorithm).CompressAsync(serializedActivityState, cancellationToken) : default; dbContext.Entry(entity).Property("SerializedActivityState").CurrentValue = compressedSerializedActivityState; dbContext.Entry(entity).Property("SerializedActivityStateCompressionAlgorithm").CurrentValue = compressionAlgorithm; @@ -102,7 +102,7 @@ public class EFCoreActivityExecutionStore : IActivityExecutionStore if (!string.IsNullOrWhiteSpace(json)) { var compressionAlgorithm = (string?)dbContext.Entry(entity).Property("SerializedActivityStateCompressionAlgorithm").CurrentValue ?? nameof(None); - var compressionStrategy = _compressionStrategyResolver.Resolve(compressionAlgorithm); + var compressionStrategy = _compressionCodecResolver.Resolve(compressionAlgorithm); json = await compressionStrategy.DecompressAsync(json, cancellationToken); var dictionary = JsonSerializer.Deserialize>(json); return dictionary?.ToDictionary(x => x.Key, x => (object)x.Value); diff --git a/src/modules/Elsa.Workflows.Management/Compression/GZip.cs b/src/modules/Elsa.Workflows.Management/Compression/GZip.cs index 3e45b7257..90da46cc3 100644 --- a/src/modules/Elsa.Workflows.Management/Compression/GZip.cs +++ b/src/modules/Elsa.Workflows.Management/Compression/GZip.cs @@ -7,7 +7,7 @@ namespace Elsa.Workflows.Management.Compression; /// /// Represents a GZip compression strategy. /// -public class GZip : ICompressionStrategy +public class GZip : ICompressionCodec { /// public async ValueTask CompressAsync(string input, CancellationToken cancellationToken) diff --git a/src/modules/Elsa.Workflows.Management/Compression/None.cs b/src/modules/Elsa.Workflows.Management/Compression/None.cs index 3072c74b3..d06420f63 100644 --- a/src/modules/Elsa.Workflows.Management/Compression/None.cs +++ b/src/modules/Elsa.Workflows.Management/Compression/None.cs @@ -5,7 +5,7 @@ namespace Elsa.Workflows.Management.Compression; /// /// Represents a compression strategy that does not compress or decompress the input. /// -public class None : ICompressionStrategy +public class None : ICompressionCodec { /// public ValueTask CompressAsync(string input, CancellationToken cancellationToken) diff --git a/src/modules/Elsa.Workflows.Management/Compression/Zstd.cs b/src/modules/Elsa.Workflows.Management/Compression/Zstd.cs new file mode 100644 index 000000000..448ddeb18 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Compression/Zstd.cs @@ -0,0 +1,37 @@ +using System.Text; +using Elsa.Workflows.Management.Contracts; +using IronCompress; + +namespace Elsa.Workflows.Management.Compression; + +/// +/// Represents a ZSTD compression strategy. +/// +public class Zstd : ICompressionCodec +{ + private Iron Iron { get; set; } = new(); + + /// + public ValueTask CompressAsync(string input, CancellationToken cancellationToken = default) + { + var inputBytes = Encoding.UTF8.GetBytes(input); + var span = inputBytes.AsSpan(); + var result = Iron.Compress(Codec.Zstd, span); + var compressedBytes = result.AsSpan(); + var compressedString = Convert.ToBase64String(compressedBytes); + + return new (compressedString); + } + + /// + public ValueTask DecompressAsync(string input, CancellationToken cancellationToken = default) + { + var inputBytes = Convert.FromBase64String(input); + var span = inputBytes.AsSpan(); + var result = Iron.Decompress(Codec.Zstd, span); + var decompressedBytes = result.AsSpan(); + var decompressedString = Encoding.UTF8.GetString(decompressedBytes); + + return new (decompressedString); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Contracts/ICompressionStrategy.cs b/src/modules/Elsa.Workflows.Management/Contracts/ICompressionCodec.cs similarity index 92% rename from src/modules/Elsa.Workflows.Management/Contracts/ICompressionStrategy.cs rename to src/modules/Elsa.Workflows.Management/Contracts/ICompressionCodec.cs index d000ec2d7..f8bedefa1 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/ICompressionStrategy.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/ICompressionCodec.cs @@ -3,7 +3,7 @@ namespace Elsa.Workflows.Management.Contracts; /// /// Represents a compression strategy. /// -public interface ICompressionStrategy +public interface ICompressionCodec { /// /// Compresses the input. diff --git a/src/modules/Elsa.Workflows.Management/Contracts/ICompressionCodecResolver.cs b/src/modules/Elsa.Workflows.Management/Contracts/ICompressionCodecResolver.cs new file mode 100644 index 000000000..cb7d234fc --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Contracts/ICompressionCodecResolver.cs @@ -0,0 +1,12 @@ +namespace Elsa.Workflows.Management.Contracts; + +/// +/// Resolves a from its name. +/// +public interface ICompressionCodecResolver +{ + /// + /// Resolves a from its name. + /// + ICompressionCodec Resolve(string name); +} \ 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 deleted file mode 100644 index df7824b3c..000000000 --- a/src/modules/Elsa.Workflows.Management/Contracts/ICompressionStrategyResolver.cs +++ /dev/null @@ -1,12 +0,0 @@ -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/Elsa.Workflows.Management.csproj b/src/modules/Elsa.Workflows.Management/Elsa.Workflows.Management.csproj index b539f922b..b0cc890fe 100644 --- a/src/modules/Elsa.Workflows.Management/Elsa.Workflows.Management.csproj +++ b/src/modules/Elsa.Workflows.Management/Elsa.Workflows.Management.csproj @@ -9,6 +9,7 @@ + diff --git a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs index 5ce9ceb9f..df1288171 100644 --- a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs +++ b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs @@ -188,9 +188,10 @@ public class WorkflowManagementFeature : FeatureBase .AddScoped() .AddSingleton() .AddSingleton() - .AddSingleton() - .AddSingleton() - .AddSingleton() + .AddSingleton() + .AddSingleton() + .AddSingleton() + .AddSingleton() ; Services.AddNotificationHandlersFrom(GetType()); diff --git a/src/modules/Elsa.Workflows.Management/Services/CompressionCodecResolver.cs b/src/modules/Elsa.Workflows.Management/Services/CompressionCodecResolver.cs new file mode 100644 index 000000000..74c0e1501 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Services/CompressionCodecResolver.cs @@ -0,0 +1,24 @@ +using Elsa.Workflows.Management.Compression; +using Elsa.Workflows.Management.Contracts; + +namespace Elsa.Workflows.Management.Services; + +/// +public class CompressionCodecResolver : ICompressionCodecResolver +{ + private readonly IDictionary _codecs; + + /// + /// Initializes a new instance of the class. + /// + public CompressionCodecResolver(IEnumerable codecs) + { + _codecs = codecs.ToDictionary(c => c.GetType().Name); + } + + /// + public ICompressionCodec Resolve(string name) + { + return _codecs.TryGetValue(name, out var codec) ? codec : new None(); + } +} \ 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 deleted file mode 100644 index bb435aa75..000000000 --- a/src/modules/Elsa.Workflows.Management/Services/CompressionStrategyResolver.cs +++ /dev/null @@ -1,14 +0,0 @@ -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