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.
This commit is contained in:
Sipke Schoorstra 2024-02-04 12:18:55 +01:00
parent 0bc46d3f74
commit 611c72fa26
14 changed files with 96 additions and 46 deletions

View file

@ -1,4 +1,4 @@
<wpf:ResourceDictionary xml:space="preserve" xmlns:x="http://schemas.microsoft.com/winfx/2006/xaml" xmlns:s="clr-namespace:System;assembly=mscorlib" xmlns:ss="urn:shemas-jetbrains-com:settings-storage-xaml" xmlns:wpf="http://schemas.microsoft.com/winfx/2006/xaml/presentation">
<wpf:ResourceDictionary xml:space="preserve" xmlns:x="http://schemas.microsoft.com/winfx/2006/xaml" xmlns:s="clr-namespace:System;assembly=mscorlib" xmlns:ss="urn:shemas-jetbrains-com:settings-storage-xaml" xmlns:wpf="http://schemas.microsoft.com/winfx/2006/xaml/presentation">
<s:Boolean x:Key="/Default/CodeStyle/CodeFormatting/CSharpFormat/KEEP_EXISTING_INITIALIZER_ARRANGEMENT/@EntryValue">False</s:Boolean>
<s:Int64 x:Key="/Default/CodeStyle/CodeFormatting/CSharpFormat/MAX_ARRAY_INITIALIZER_ELEMENTS_ON_LINE/@EntryValue">1000</s:Int64>
<s:Int64 x:Key="/Default/CodeStyle/CodeFormatting/CSharpFormat/MAX_INITIALIZER_ELEMENTS_ON_LINE/@EntryValue">1</s:Int64>
@ -17,4 +17,5 @@
<s:Boolean x:Key="/Default/UserDictionary/Words/=Postgre/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=startable/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=Telnyx/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=Unschedule/@EntryIndexedValue">True</s:Boolean></wpf:ResourceDictionary>
<s:Boolean x:Key="/Default/UserDictionary/Words/=Unschedule/@EntryIndexedValue">True</s:Boolean>
<s:Boolean x:Key="/Default/UserDictionary/Words/=Zstd/@EntryIndexedValue">True</s:Boolean></wpf:ResourceDictionary>

View file

@ -110,7 +110,7 @@ services
});
if(useZipCompression)
management.SetCompressionAlgorithm(nameof(GZip));
management.SetCompressionAlgorithm(nameof(Zstd));
})
.UseWorkflowRuntime(runtime =>
{

View file

@ -23,7 +23,7 @@ public class EFCoreWorkflowInstanceStore : IWorkflowInstanceStore
{
private readonly EntityStore<ManagementElsaDbContext, WorkflowInstance> _store;
private readonly IWorkflowStateSerializer _workflowStateSerializer;
private readonly ICompressionStrategyResolver _compressionStrategyResolver;
private readonly ICompressionCodecResolver _compressionCodecResolver;
private readonly IOptions<ManagementOptions> _options;
/// <summary>
@ -32,12 +32,12 @@ public class EFCoreWorkflowInstanceStore : IWorkflowInstanceStore
public EFCoreWorkflowInstanceStore(
EntityStore<ManagementElsaDbContext, WorkflowInstance> store,
IWorkflowStateSerializer workflowStateSerializer,
ICompressionStrategyResolver compressionStrategyResolver,
ICompressionCodecResolver compressionCodecResolver,
IOptions<ManagementOptions> 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))
{

View file

@ -25,7 +25,7 @@ public class EFCoreActivityExecutionStore : IActivityExecutionStore
private readonly EntityStore<RuntimeElsaDbContext, ActivityExecutionRecord> _store;
private readonly ISafeSerializer _safeSerializer;
private readonly IPayloadSerializer _payloadSerializer;
private readonly ICompressionStrategyResolver _compressionStrategyResolver;
private readonly ICompressionCodecResolver _compressionCodecResolver;
private readonly IOptions<ManagementOptions> _options;
/// <summary>
@ -35,13 +35,13 @@ public class EFCoreActivityExecutionStore : IActivityExecutionStore
EntityStore<RuntimeElsaDbContext, ActivityExecutionRecord> store,
ISafeSerializer safeSerializer,
IPayloadSerializer payloadSerializer,
ICompressionStrategyResolver compressionStrategyResolver,
ICompressionCodecResolver compressionCodecResolver,
IOptions<ManagementOptions> 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<IDictionary<string, object>>(json);
return dictionary?.ToDictionary(x => x.Key, x => (object)x.Value);

View file

@ -7,7 +7,7 @@ namespace Elsa.Workflows.Management.Compression;
/// <summary>
/// Represents a GZip compression strategy.
/// </summary>
public class GZip : ICompressionStrategy
public class GZip : ICompressionCodec
{
/// <inheritdoc />
public async ValueTask<string> CompressAsync(string input, CancellationToken cancellationToken)

View file

@ -5,7 +5,7 @@ namespace Elsa.Workflows.Management.Compression;
/// <summary>
/// Represents a compression strategy that does not compress or decompress the input.
/// </summary>
public class None : ICompressionStrategy
public class None : ICompressionCodec
{
/// <inheritdoc />
public ValueTask<string> CompressAsync(string input, CancellationToken cancellationToken)

View file

@ -0,0 +1,37 @@
using System.Text;
using Elsa.Workflows.Management.Contracts;
using IronCompress;
namespace Elsa.Workflows.Management.Compression;
/// <summary>
/// Represents a ZSTD compression strategy.
/// </summary>
public class Zstd : ICompressionCodec
{
private Iron Iron { get; set; } = new();
/// <inheritdoc />
public ValueTask<string> 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);
}
/// <inheritdoc />
public ValueTask<string> 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);
}
}

View file

@ -3,7 +3,7 @@ namespace Elsa.Workflows.Management.Contracts;
/// <summary>
/// Represents a compression strategy.
/// </summary>
public interface ICompressionStrategy
public interface ICompressionCodec
{
/// <summary>
/// Compresses the input.

View file

@ -0,0 +1,12 @@
namespace Elsa.Workflows.Management.Contracts;
/// <summary>
/// Resolves a <see cref="ICompressionCodec"/> from its name.
/// </summary>
public interface ICompressionCodecResolver
{
/// <summary>
/// Resolves a <see cref="ICompressionCodec"/> from its name.
/// </summary>
ICompressionCodec Resolve(string name);
}

View file

@ -1,12 +0,0 @@
namespace Elsa.Workflows.Management.Contracts;
/// <summary>
/// Resolves a <see cref="ICompressionStrategy"/> from its name.
/// </summary>
public interface ICompressionStrategyResolver
{
/// <summary>
/// Resolves a <see cref="ICompressionStrategy"/> from its name.
/// </summary>
ICompressionStrategy Resolve(string name);
}

View file

@ -9,6 +9,7 @@
<ItemGroup>
<PackageReference Include="Humanizer.Core" Version="2.14.1"/>
<PackageReference Include="IronCompress" Version="1.5.1" />
<PackageReference Include="System.Linq.Async" Version="6.0.1"/>
</ItemGroup>

View file

@ -188,9 +188,10 @@ public class WorkflowManagementFeature : FeatureBase
.AddScoped<WorkflowDefinitionMapper>()
.AddSingleton<VariableDefinitionMapper>()
.AddSingleton<WorkflowStateMapper>()
.AddSingleton<ICompressionStrategyResolver, CompressionStrategyResolver>()
.AddSingleton<ICompressionStrategy, None>()
.AddSingleton<ICompressionStrategy, GZip>()
.AddSingleton<ICompressionCodecResolver, CompressionCodecResolver>()
.AddSingleton<ICompressionCodec, None>()
.AddSingleton<ICompressionCodec, GZip>()
.AddSingleton<ICompressionCodec, Zstd>()
;
Services.AddNotificationHandlersFrom(GetType());

View file

@ -0,0 +1,24 @@
using Elsa.Workflows.Management.Compression;
using Elsa.Workflows.Management.Contracts;
namespace Elsa.Workflows.Management.Services;
/// <inheritdoc />
public class CompressionCodecResolver : ICompressionCodecResolver
{
private readonly IDictionary<string, ICompressionCodec> _codecs;
/// <summary>
/// Initializes a new instance of the <see cref="CompressionCodecResolver"/> class.
/// </summary>
public CompressionCodecResolver(IEnumerable<ICompressionCodec> codecs)
{
_codecs = codecs.ToDictionary(c => c.GetType().Name);
}
/// <inheritdoc />
public ICompressionCodec Resolve(string name)
{
return _codecs.TryGetValue(name, out var codec) ? codec : new None();
}
}

View file

@ -1,14 +0,0 @@
using Elsa.Workflows.Management.Compression;
using Elsa.Workflows.Management.Contracts;
namespace Elsa.Workflows.Management.Services;
/// <inheritdoc />
public class CompressionStrategyResolver(IEnumerable<ICompressionStrategy> strategies) : ICompressionStrategyResolver
{
/// <inheritdoc />
public ICompressionStrategy Resolve(string name)
{
return strategies.FirstOrDefault(s => s.GetType().Name == name) ?? new None();
}
}