diff --git a/Directory.Packages.props b/Directory.Packages.props index f755c1f93..684a81fa2 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -180,6 +180,7 @@ 3.5.0-preview.1092 + 9.0.7 @@ -209,7 +210,7 @@ - + @@ -234,7 +235,7 @@ - + @@ -293,8 +294,8 @@ - - + + @@ -302,36 +303,36 @@ - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + diff --git a/src/modules/Elsa.Common/Entities/Entity.cs b/src/modules/Elsa.Common/Entities/Entity.cs index 2a3d2cd6e..3da8462f6 100644 --- a/src/modules/Elsa.Common/Entities/Entity.cs +++ b/src/modules/Elsa.Common/Entities/Entity.cs @@ -8,7 +8,7 @@ public abstract class Entity /// /// Gets or sets the ID of this entity. /// - public string Id { get; set; } = default!; + public string Id { get; set; } = null!; /// /// Gets or sets the ID of the tenant that own this entity. diff --git a/src/modules/Elsa.Expressions.JavaScript/Handlers/ConfigureEngineWithCommonFunctions.cs b/src/modules/Elsa.Expressions.JavaScript/Handlers/ConfigureEngineWithCommonFunctions.cs index 922974e49..1e85c6698 100644 --- a/src/modules/Elsa.Expressions.JavaScript/Handlers/ConfigureEngineWithCommonFunctions.cs +++ b/src/modules/Elsa.Expressions.JavaScript/Handlers/ConfigureEngineWithCommonFunctions.cs @@ -59,6 +59,8 @@ public class ConfigureEngineWithCommonFunctions(IOptions options) : engine.SetValue("bytesFromBase64", (Func)(value => Convert.FromBase64String(value))); engine.SetValue("stringToBase64", (Func)(value => Convert.ToBase64String(Encoding.UTF8.GetBytes(value)))); engine.SetValue("stringFromBase64", (Func)(value => Encoding.UTF8.GetString(Convert.FromBase64String(value)))); + engine.SetValue("streamToBytes", (Func)(value => StreamToBytes(value))); + engine.SetValue("streamToBase64", (Func)(value => Convert.ToBase64String(StreamToBytes(value)))); // Deprecated, use newGuidString instead. engine.SetValue("getGuidString", (Func)(() => Guid.NewGuid().ToString())); @@ -67,7 +69,7 @@ public class ConfigureEngineWithCommonFunctions(IOptions options) : engine.SetValue("getShortGuid", (Func)(() => Regex.Replace(Convert.ToBase64String(Guid.NewGuid().ToByteArray()), "[/+=]", ""))); return Task.CompletedTask; } - + private string Serialize(object value) { return JsonSerializer.Serialize(value, _jsonSerializerOptions); @@ -82,4 +84,11 @@ public class ConfigureEngineWithCommonFunctions(IOptions options) : options.Converters.Add(new JsonStringEnumConverter()); return options; } + + private byte[] StreamToBytes(Stream stream) + { + using var memoryStream = new MemoryStream(); + stream.CopyTo(memoryStream); + return memoryStream.ToArray(); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Expressions.JavaScript/Providers/CommonFunctionsDefinitionProvider.cs b/src/modules/Elsa.Expressions.JavaScript/Providers/CommonFunctionsDefinitionProvider.cs index b612906e5..de943b76e 100644 --- a/src/modules/Elsa.Expressions.JavaScript/Providers/CommonFunctionsDefinitionProvider.cs +++ b/src/modules/Elsa.Expressions.JavaScript/Providers/CommonFunctionsDefinitionProvider.cs @@ -136,6 +136,16 @@ internal class CommonFunctionsDefinitionProvider(ITypeAliasRegistry typeAliasReg .Name("stringToBase64") .Parameter("value", "string") .ReturnType("string")); + + yield return CreateFunctionDefinition(builder => builder + .Name("streamToBytes") + .Parameter("value", "Stream") + .ReturnType("Byte[]")); + + yield return CreateFunctionDefinition(builder => builder + .Name("streamToBase64") + .Parameter("value", "Stream") + .ReturnType("string")); if (!options.Value.DisableWrappers) { @@ -151,7 +161,7 @@ internal class CommonFunctionsDefinitionProvider(ITypeAliasRegistry typeAliasReg // set{Variable}. yield return CreateFunctionDefinition(builder => builder.Name($"set{pascalName}").Parameter("value", typeAlias)); - } + } } } } \ No newline at end of file diff --git a/src/modules/Elsa.IO.Compression/Activities/CreateZipArchive.cs b/src/modules/Elsa.IO.Compression/Activities/CreateZipArchive.cs index c7eee4d13..e322d7594 100644 --- a/src/modules/Elsa.IO.Compression/Activities/CreateZipArchive.cs +++ b/src/modules/Elsa.IO.Compression/Activities/CreateZipArchive.cs @@ -80,7 +80,7 @@ public class CreateZipArchive : CodeActivity try { - using var zipArchive = new ZipArchive(zipStream, ZipArchiveMode.Create, leaveOpen: true); + using var zipArchive = new ZipArchive(zipStream, ZipArchiveMode.Update, leaveOpen: true); var entryIndex = 0; var compressionLevel = CompressionLevel.Get(context); @@ -122,19 +122,87 @@ public class CreateZipArchive : CodeActivity CompressionLevel compressionLevel) { var binaryContent = await resolver.ResolveAsync(entryContent, context.CancellationToken); - - var entryName = binaryContent.Name?.GetNameAndExtension() + + var entryName = binaryContent.Name?.GetNameAndExtension() ?? string.Format(DefaultEntryNameFormat, entryIndex + 1); + // Get a unique name following Windows convention + entryName = GetUniqueEntryName(zipArchive, entryName); + var archiveEntry = zipArchive.CreateEntry(entryName, compressionLevel); await using var entryStream = archiveEntry.Open(); await binaryContent.Stream.CopyToAsync(entryStream, context.CancellationToken); await entryStream.FlushAsync(context.CancellationToken); - + if (entryContent is not Stream) { await binaryContent.Stream.DisposeAsync(); } } + + private static string GetUniqueEntryName(ZipArchive zipArchive, string originalName) + { + var filenameWithoutExtension = Path.GetFileNameWithoutExtension(originalName); + var extension = Path.GetExtension(originalName); + + var originalExists = false; + var highestIndex = 0; + + foreach (var entry in zipArchive.Entries) + { + if (!entry.Name.Equals(originalName, StringComparison.OrdinalIgnoreCase)) + { + continue; + } + + originalExists = true; + + var entryNameWithoutExtension = Path.GetFileNameWithoutExtension(entry.Name); + var entryExtension = Path.GetExtension(entry.Name); + + // Only process entries with the same extension + if (!entryExtension.Equals(extension, StringComparison.OrdinalIgnoreCase)) + continue; + + // Check if this entry follows our naming pattern + highestIndex = HighestEntryNameIndex(entryNameWithoutExtension, filenameWithoutExtension, highestIndex); + } + + if (!originalExists) + { + return originalName; + } + + return $"{filenameWithoutExtension}({highestIndex + 1}){extension}"; + } + + private static int HighestEntryNameIndex(string entryNameWithoutExtension, string filenameWithoutExtension, + int highestIndex) + { + if (!entryNameWithoutExtension.StartsWith(filenameWithoutExtension, StringComparison.OrdinalIgnoreCase) || + entryNameWithoutExtension.Length <= filenameWithoutExtension.Length || + entryNameWithoutExtension[filenameWithoutExtension.Length] != '(') + { + return highestIndex; + } + + // Extract the number between parentheses + var closingParenIndex = entryNameWithoutExtension.LastIndexOf(')'); + if (closingParenIndex <= filenameWithoutExtension.Length + 1) + { + return highestIndex; + } + + var indexStr = entryNameWithoutExtension.Substring( + filenameWithoutExtension.Length + 1, + closingParenIndex - filenameWithoutExtension.Length - 1); + + if (int.TryParse(indexStr, out var index)) + { + highestIndex = Math.Max(highestIndex, index); + } + + return highestIndex; + } } \ No newline at end of file diff --git a/src/modules/Elsa.IO.Compression/Services/Strategies/ZipEntryContentStrategy.cs b/src/modules/Elsa.IO.Compression/Services/Strategies/ZipEntryContentStrategy.cs index 14757faa1..feaca88e0 100644 --- a/src/modules/Elsa.IO.Compression/Services/Strategies/ZipEntryContentStrategy.cs +++ b/src/modules/Elsa.IO.Compression/Services/Strategies/ZipEntryContentStrategy.cs @@ -36,7 +36,7 @@ public class ZipEntryContentStrategy(IServiceProvider serviceProvider) : IConten var innerContentName = innerContent.Name?.GetNameAndExtension(); var innerContentExtension = Path.GetExtension(innerContentName); innerContent.Name = !string.IsNullOrWhiteSpace(innerContentExtension) - ? zipEntry.EntryName + innerContentExtension + ? Path.HasExtension(zipEntry.EntryName) ? zipEntry.EntryName : zipEntry.EntryName + innerContentExtension : innerContent.Name; return innerContent; diff --git a/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs b/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs index 0bb796073..45bf158ef 100644 --- a/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs +++ b/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs @@ -85,6 +85,62 @@ public static class ContentTypeExtensions return Path.GetExtension(filePath).ToLowerInvariant(); } + public static bool IsBase64String(this string s) + { + if (string.IsNullOrWhiteSpace(s)) + return false; + + s = s.Trim(); + + // Length must be divisible by 4 + if (s.Length % 4 != 0) + return false; + + // Check padding position and count + var paddingIndex = s.IndexOf('='); + + switch (paddingIndex) + { + // Padding cannot be at index 0 + case 0: + // Padding must be at the end + case > 0 when paddingIndex < s.Length - 2: + // All characters after first '=' must also be '=' + case > 0 when s[paddingIndex..].Any(c => c != '='): + return false; + } + + // Check for valid Base64 characters + for (var i = 0; i < paddingIndex; i++) + { + var c = s[i]; + var isValid = + c is >= 'A' and <= 'Z' || + c is >= 'a' and <= 'z' || + c is >= '0' and <= '9' || + c == '+' || c == '/'; + + if (!isValid) + return false; + } + + // Additional check for short strings that are just lowercase+numbers + // This catches "whatever" and similar false positives + if (s.Length <= 10 && s.All(c => char.IsLower(c) || char.IsDigit(c))) + return false; + + // Try actual decoding + try + { + _ = Convert.FromBase64String(s); + return true; + } + catch + { + return false; + } + } + private static string DetermineExtensionFromMimeType(string mimeType) { if (mimeType.Contains("/pdf")) diff --git a/src/modules/Elsa.IO/Services/Strategies/Base64ContentStrategy.cs b/src/modules/Elsa.IO/Services/Strategies/Base64ContentStrategy.cs index a1d621c9b..61e7bb39d 100644 --- a/src/modules/Elsa.IO/Services/Strategies/Base64ContentStrategy.cs +++ b/src/modules/Elsa.IO/Services/Strategies/Base64ContentStrategy.cs @@ -51,13 +51,8 @@ public class Base64ContentStrategy : IContentResolverStrategy private static bool IsBase64String(string base64) { - if (IsUriDataBase64String(base64)) - { - return true; - } - - var buffer = new Span(new byte[base64.Length]); - return Convert.TryFromBase64String(base64, buffer , out _); + return IsUriDataBase64String(base64) + || base64.IsBase64String(); } private static bool IsUriDataBase64String(string base64) diff --git a/src/modules/Elsa.Workflows.Core/VariableStorageDrivers/WorkflowInstanceStorageDriver.cs b/src/modules/Elsa.Workflows.Core/VariableStorageDrivers/WorkflowInstanceStorageDriver.cs index 566473c3a..bb93720e3 100644 --- a/src/modules/Elsa.Workflows.Core/VariableStorageDrivers/WorkflowInstanceStorageDriver.cs +++ b/src/modules/Elsa.Workflows.Core/VariableStorageDrivers/WorkflowInstanceStorageDriver.cs @@ -4,6 +4,7 @@ using System.Text.Json.Nodes; using Elsa.Expressions.Helpers; using Elsa.Extensions; using JetBrains.Annotations; +using Microsoft.Extensions.Logging; namespace Elsa.Workflows; @@ -12,7 +13,7 @@ namespace Elsa.Workflows; /// [Display(Name = "Workflow Instance")] [UsedImplicitly] -public class WorkflowInstanceStorageDriver(IPayloadSerializer payloadSerializer) : IStorageDriver +public class WorkflowInstanceStorageDriver(IPayloadSerializer payloadSerializer, ILogger logger) : IStorageDriver { /// /// The key used to store the variables in the workflow state. @@ -29,8 +30,18 @@ public class WorkflowInstanceStorageDriver(IPayloadSerializer payloadSerializer) { UpdateVariablesDictionary(context, dictionary => { - var node = JsonSerializer.SerializeToNode(value); - dictionary[id] = node; + try + { + var node = JsonSerializer.SerializeToNode(value); + dictionary[id] = node; + } + catch (Exception ex) when (ex is JsonException or NotSupportedException or ObjectDisposedException) + { + logger.LogWarning(ex, "Failed to serialize variable '{VariableId}' of type '{VariableType}' for workflow instance storage. The variable will be skipped.", + id, value?.GetType().FullName ?? "null"); + + dictionary.Remove(id); + } }); return ValueTask.CompletedTask; } diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IActivityExecutonMapper.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IActivityExecutonMapper.cs index 05b943522..c533c2b00 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IActivityExecutonMapper.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IActivityExecutonMapper.cs @@ -10,11 +10,5 @@ public interface IActivityExecutionMapper /// /// Maps an activity execution context to an activity execution record. /// - ActivityExecutionRecord Map(ActivityExecutionContext source); - - /// - /// Maps an activity execution context to an activity execution record. - /// - [Obsolete( "Use Map instead.", error: false)] Task MapAsync(ActivityExecutionContext source); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs b/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs index bf47b6a24..2b6bd1e74 100644 --- a/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs +++ b/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs @@ -1,3 +1,5 @@ +using System.ComponentModel.DataAnnotations.Schema; +using System.Text.Json.Serialization; using Elsa.Common; using Elsa.Common.Entities; using Elsa.Workflows.State; @@ -7,18 +9,18 @@ namespace Elsa.Workflows.Runtime.Entities; /// /// Represents a single activity execution of an activity instance. /// -public class ActivityExecutionRecord : Entity, ILogRecord +public partial class ActivityExecutionRecord : Entity, ILogRecord { /// /// Gets or sets the workflow instance ID. /// public string WorkflowInstanceId { get; set; } = null!; - + /// /// Gets or sets the activity ID. /// public string ActivityId { get; set; } = null!; - + /// /// Gets or sets the activity node ID. /// @@ -38,17 +40,17 @@ public class ActivityExecutionRecord : Entity, ILogRecord /// The name of the activity. /// public string? ActivityName { get; set; } - + /// /// The state of the activity at the time this record is created or last updated. /// public IDictionary? ActivityState { get; set; } - + /// /// Any additional payload associated with the log record. /// public IDictionary? Payload { get; set; } - + /// /// Any outputs provided by the activity. /// @@ -58,7 +60,7 @@ public class ActivityExecutionRecord : Entity, ILogRecord /// Any properties provided by the activity. /// public IDictionary? Properties { get; set; } - + /// /// Lightweight metadata associated with the activity execution. /// This information will be retained as part of the activity execution summary record. @@ -79,7 +81,7 @@ public class ActivityExecutionRecord : Entity, ILogRecord /// Gets or sets whether the activity has any bookmarks. /// public bool HasBookmarks { get; set; } - + /// /// Gets or sets the status of the activity. /// @@ -89,9 +91,14 @@ public class ActivityExecutionRecord : Entity, ILogRecord /// Gets or sets the aggregated count of faults encountered during the execution of the activity instance and its descendants. /// public int AggregateFaultCount { get; set; } - + /// /// Gets or sets the time at which the activity execution completed. /// public DateTimeOffset? CompletedAt { get; set; } +} + +public partial class ActivityExecutionRecord +{ + [NotMapped] [JsonIgnore] public ActivityExecutionRecordSnapshot? SerializedSnapshot { get; set; } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextExtensions.cs b/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextEventExtensions.cs similarity index 96% rename from src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextExtensions.cs rename to src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextEventExtensions.cs index e222e8dbc..151a66e2d 100644 --- a/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextExtensions.cs +++ b/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextEventExtensions.cs @@ -7,7 +7,7 @@ using Elsa.Workflows.Runtime.Stimuli; // ReSharper disable once CheckNamespace namespace Elsa.Extensions; -public static class ActivityExecutionContextExtensions +public static class ActivityExecutionContextEventExtensions { /// /// Suspends the current activity's execution and waits for a specified event to occur before continuing. diff --git a/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs b/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs new file mode 100644 index 000000000..c2706cc1e --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs @@ -0,0 +1,32 @@ +using Elsa.Workflows; +using Elsa.Workflows.Runtime; +using Elsa.Workflows.Runtime.Entities; + +// ReSharper disable once CheckNamespace +namespace Elsa.Extensions; + +public static class ActivityExecutionContextRecordExtensions +{ + private const string ActivityExecutionRecordKey = "CapturedActivityExecutionRecord"; + + public static async Task CaptureActivityExecutionRecordAsync(this ActivityExecutionContext context) + { + var mapper = context.GetRequiredService(); + var record = await mapper.MapAsync(context); + context.TransientProperties[ActivityExecutionRecordKey] = record; + } + + public static ActivityExecutionRecord? GetCapturedActivityExecutionRecord(this ActivityExecutionContext context) + { + return context.TransientProperties.TryGetValue(ActivityExecutionRecordKey, out var record) ? (ActivityExecutionRecord?)record : null; + } + + public static async Task GetOrMapCapturedActivityExecutionRecordAsync(this ActivityExecutionContext context) + { + if(context.TransientProperties.TryGetValue(ActivityExecutionRecordKey, out var record)) + return (ActivityExecutionRecord)record; + + var mapper = context.GetRequiredService(); + return await mapper.MapAsync(context); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Extensions/PipelineWorkflowsFeatureExtensions.cs b/src/modules/Elsa.Workflows.Runtime/Extensions/PipelineWorkflowsFeatureExtensions.cs index 642e449eb..0224b0fb4 100644 --- a/src/modules/Elsa.Workflows.Runtime/Extensions/PipelineWorkflowsFeatureExtensions.cs +++ b/src/modules/Elsa.Workflows.Runtime/Extensions/PipelineWorkflowsFeatureExtensions.cs @@ -34,6 +34,7 @@ public static class PipelineWorkflowsFeatureExtensions .UseExecutionLogging() .UseNotifications() .UseLogPersistenceModeEvaluation() + .UseActivityExecutionLogCapturing() .UseBackgroundActivityInvoker(); configurePipeline?.Invoke(pipeline); diff --git a/src/modules/Elsa.Workflows.Runtime/LogPersistence/Extensions/ActivityExecutionPipelineBuilderExtensions.cs b/src/modules/Elsa.Workflows.Runtime/LogPersistence/Extensions/ActivityExecutionPipelineBuilderExtensions.cs index 988b6d87b..b57f15df0 100644 --- a/src/modules/Elsa.Workflows.Runtime/LogPersistence/Extensions/ActivityExecutionPipelineBuilderExtensions.cs +++ b/src/modules/Elsa.Workflows.Runtime/LogPersistence/Extensions/ActivityExecutionPipelineBuilderExtensions.cs @@ -20,4 +20,12 @@ public static class ActivityExecutionPipelineBuilderExtensions /// Installs the which evaluates log persistence modes during activity execution. /// public static IActivityExecutionPipelineBuilder UseLogPersistenceModeEvaluation(this IActivityExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder.UseMiddleware(); + + /// + /// Installs the into the activity execution pipeline to capture and map activity execution details. + /// + public static IActivityExecutionPipelineBuilder UseActivityExecutionLogCapturing(this IActivityExecutionPipelineBuilder pipelineBuilder) + { + return pipelineBuilder.UseMiddleware(); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Middleware/Activities/CaptureActivityExecutionRecordMiddleware.cs b/src/modules/Elsa.Workflows.Runtime/Middleware/Activities/CaptureActivityExecutionRecordMiddleware.cs new file mode 100644 index 000000000..986a05ef7 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Middleware/Activities/CaptureActivityExecutionRecordMiddleware.cs @@ -0,0 +1,13 @@ +using Elsa.Extensions; +using Elsa.Workflows.Pipelines.ActivityExecution; + +namespace Elsa.Workflows.Runtime.Middleware.Activities; + +public class CaptureActivityExecutionRecordMiddleware(ActivityMiddlewareDelegate next) : IActivityExecutionMiddleware +{ + public async ValueTask InvokeAsync(ActivityExecutionContext context) + { + await next(context); + await context.CaptureActivityExecutionRecordAsync(); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Models/ActivityExecutionRecordSnapshot.cs b/src/modules/Elsa.Workflows.Runtime/Models/ActivityExecutionRecordSnapshot.cs new file mode 100644 index 000000000..377850153 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Models/ActivityExecutionRecordSnapshot.cs @@ -0,0 +1,25 @@ +namespace Elsa.Workflows.Runtime; + +public class ActivityExecutionRecordSnapshot +{ + public string Id { get; set; } = null!; + public string? TenantId { get; set; } + public string WorkflowInstanceId { get; set; } = null!; + public string ActivityId { get; set; } = null!; + public string ActivityNodeId { get; set; } = null!; + public string ActivityType { get; set; } = null!; + public int ActivityTypeVersion { get; set; } + public string? ActivityName { get; set; } + public DateTimeOffset StartedAt { get; set; } + public bool HasBookmarks { get; set; } + public ActivityStatus Status { get; set; } + public int AggregateFaultCount { get; set; } + public DateTimeOffset? CompletedAt { get; set; } + public string? SerializedActivityState { get; set; } + public string? SerializedOutputs { get; set; } + public string? SerializedProperties { get; set; } + public string? SerializedPayload { get; set; } + public string? SerializedMetadata { get; set; } + public string? SerializedException { get; set; } + public string? SerializedActivityStateCompressionAlgorithm { get; set; } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs index e862e6be8..223cd55f3 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs @@ -1,13 +1,23 @@ +using Elsa.Common; +using Elsa.Common.Codecs; using Elsa.Workflows.LogPersistence; +using Elsa.Workflows.Management.Options; using Elsa.Workflows.Runtime.Entities; +using Elsa.Workflows.Runtime.Extensions; using Elsa.Workflows.State; +using Microsoft.Extensions.Options; namespace Elsa.Workflows.Runtime; /// -public class DefaultActivityExecutionMapper() : IActivityExecutionMapper +public class DefaultActivityExecutionMapper( + ISafeSerializer safeSerializer, + IPayloadSerializer payloadSerializer, + ICompressionCodecResolver compressionCodecResolver, + IOptions options) : IActivityExecutionMapper { - public ActivityExecutionRecord Map(ActivityExecutionContext source) + /// + public async Task MapAsync(ActivityExecutionContext source) { var outputs = source.GetOutputs(); var inputs = source.GetInputs(); @@ -16,8 +26,9 @@ public class DefaultActivityExecutionMapper() : IActivityExecutionMapper var persistableOutputs = GetPersistableInputOutput(outputs, persistenceMap.Outputs); var persistableProperties = GetPersistableDictionary(source.Properties!, persistenceMap.InternalState); var persistableJournalData = GetPersistableDictionary(source.JournalData!, persistenceMap.InternalState); + var cancellationToken = source.CancellationToken; - return new() + var record = new ActivityExecutionRecord { Id = source.Id, ActivityId = source.Activity.Id, @@ -38,15 +49,41 @@ public class DefaultActivityExecutionMapper() : IActivityExecutionMapper AggregateFaultCount = source.AggregateFaultCount, CompletedAt = source.CompletedAt }; + + record = record.SanitizeLogMessage(); + var compressionAlgorithm = options.Value.CompressionAlgorithm ?? nameof(None); + var serializedActivityState = record.ActivityState?.Count > 0 ? safeSerializer.Serialize(record.ActivityState) : null; + var compressedSerializedActivityState = serializedActivityState != null ? await compressionCodecResolver.Resolve(compressionAlgorithm).CompressAsync(serializedActivityState, cancellationToken) : null; + var serializedProperties = record.Properties != null ? payloadSerializer.Serialize(record.Properties) : null; + var serializedMetadata = record.Metadata != null ? payloadSerializer.Serialize(record.Metadata) : null; + record.SerializedSnapshot = new() + { + Id = record.Id, + TenantId = record.TenantId, + WorkflowInstanceId = record.WorkflowInstanceId, + ActivityId = record.ActivityId, + ActivityNodeId = record.ActivityNodeId, + ActivityType = record.ActivityType, + ActivityTypeVersion = record.ActivityTypeVersion, + ActivityName = record.ActivityName, + StartedAt = record.StartedAt, + HasBookmarks = record.HasBookmarks, + Status = record.Status, + AggregateFaultCount = record.AggregateFaultCount, + CompletedAt = record.CompletedAt, + SerializedActivityState = compressedSerializedActivityState, + SerializedActivityStateCompressionAlgorithm = compressionAlgorithm, + SerializedOutputs = record.Outputs?.Any() == true ? safeSerializer.Serialize(record.Outputs) : null, + SerializedProperties = serializedProperties, + SerializedMetadata = serializedMetadata, + SerializedException = record.Exception != null ? payloadSerializer.Serialize(record.Exception) : null, + SerializedPayload = record.Payload?.Any() == true ? payloadSerializer.Serialize(record.Payload) : null + }; + + return record; } - /// - public Task MapAsync(ActivityExecutionContext source) - { - return Task.FromResult(Map(source)); - } - - private IDictionary GetPersistableInputOutput(IDictionary state, IDictionary map) + private IDictionary GetPersistableInputOutput(IDictionary state, IDictionary map, bool deepCopy = false) { var result = new Dictionary(); foreach (var stateEntry in state) diff --git a/src/modules/Elsa.Workflows.Runtime/Services/StoreActivityExecutionLogSink.cs b/src/modules/Elsa.Workflows.Runtime/Services/StoreActivityExecutionLogSink.cs index 01217cda1..72325955b 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/StoreActivityExecutionLogSink.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/StoreActivityExecutionLogSink.cs @@ -1,3 +1,4 @@ +using Elsa.Extensions; using Elsa.Mediator.Contracts; using Elsa.Workflows.Runtime.Entities; using Elsa.Workflows.Runtime.Notifications; @@ -22,7 +23,7 @@ public class StoreActivityExecutionLogSink( if (activityExecutionContexts.Count == 0) return; - var records = activityExecutionContexts.Select(mapper.Map).ToList(); + var records = await Task.WhenAll(activityExecutionContexts.Select(x => x.GetOrMapCapturedActivityExecutionRecordAsync())); await activityExecutionStore.SaveManyAsync(records, cancellationToken); // Untaint activity execution contexts.