From 1f0e67926f9c285bd9c93a66a163112222780b36 Mon Sep 17 00:00:00 2001 From: "lucas.hipolito" Date: Tue, 8 Jul 2025 11:26:06 +0200 Subject: [PATCH 01/14] Fixing wrong strategy selection for some string content --- .../Extensions/ContentTypeExtensions.cs | 39 +++++++++++++++++++ .../Strategies/Base64ContentStrategy.cs | 9 +---- 2 files changed, 41 insertions(+), 7 deletions(-) diff --git a/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs b/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs index 0bb796073..45e82a262 100644 --- a/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs +++ b/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs @@ -85,6 +85,45 @@ 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 valid Base64 characters + for (var i = 0; i < s.Length; i++) + { + var c = s[i]; + + var isValid = + c is >= 'A' and <= 'Z' || + c is >= 'a' and <= 'z' || + c is >= '0' and <= '9' || + c == '+' || c == '/' || c == '='; + + if (!isValid) + return false; + } + + // Try actual decoding and roundtrip + try + { + var data = Convert.FromBase64String(s); + var reEncoded = Convert.ToBase64String(data); + return s == reEncoded; + } + 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..168e13e11 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) From 9e3e93c580b3bdfb12e0b29415966649c84d1235 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Lucas=20Hip=C3=B3lito?= Date: Tue, 8 Jul 2025 20:31:35 +0200 Subject: [PATCH 02/14] changing logic to include non-Data URI scenarios changing logic to include non-Data URI scenarios Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --- .../Elsa.IO/Services/Strategies/Base64ContentStrategy.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/modules/Elsa.IO/Services/Strategies/Base64ContentStrategy.cs b/src/modules/Elsa.IO/Services/Strategies/Base64ContentStrategy.cs index 168e13e11..61e7bb39d 100644 --- a/src/modules/Elsa.IO/Services/Strategies/Base64ContentStrategy.cs +++ b/src/modules/Elsa.IO/Services/Strategies/Base64ContentStrategy.cs @@ -52,7 +52,7 @@ public class Base64ContentStrategy : IContentResolverStrategy private static bool IsBase64String(string base64) { return IsUriDataBase64String(base64) - && base64.IsBase64String(); + || base64.IsBase64String(); } private static bool IsUriDataBase64String(string base64) From 35cdaad7ef02a9a775029fbecb87cc2be827826e Mon Sep 17 00:00:00 2001 From: "lucas.hipolito" Date: Thu, 10 Jul 2025 10:07:01 +0200 Subject: [PATCH 03/14] Implementing better base64 evaluation + appropriate zip entry naming --- .../Activities/CreateZipArchive.cs | 64 +++++++++++++++++-- .../Extensions/ContentTypeExtensions.cs | 34 +++++++--- 2 files changed, 84 insertions(+), 14 deletions(-) diff --git a/src/modules/Elsa.IO.Compression/Activities/CreateZipArchive.cs b/src/modules/Elsa.IO.Compression/Activities/CreateZipArchive.cs index c7eee4d13..1e373731a 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,73 @@ 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) + { + // If no duplicate exists, use the original name + if (!zipArchive.Entries.Any(entry => entry.Name.Equals(originalName, StringComparison.OrdinalIgnoreCase))) + { + return originalName; + } + + // Split the name into filename and extension + string filenameWithoutExtension = Path.GetFileNameWithoutExtension(originalName); + string extension = Path.GetExtension(originalName); + + // Find the highest index used for this filename pattern + int highestIndex = 0; + + // Check for the original name and any name with pattern "name(n).ext" + foreach (var entry in zipArchive.Entries) + { + if (entry.Name.Equals(originalName, StringComparison.OrdinalIgnoreCase)) + continue; // Skip the exact match as we already know it exists + + string entryNameWithoutExt = Path.GetFileNameWithoutExtension(entry.Name); + string entryExt = Path.GetExtension(entry.Name); + + if (!entryExt.Equals(extension, StringComparison.OrdinalIgnoreCase)) + continue; // Different extension + + if (entryNameWithoutExt.StartsWith(filenameWithoutExtension, StringComparison.OrdinalIgnoreCase) && + entryNameWithoutExt.Length > filenameWithoutExtension.Length && + entryNameWithoutExt[filenameWithoutExtension.Length] == '(') + { + // Extract the number between parentheses + var closingParenIndex = entryNameWithoutExt.LastIndexOf(')'); + if (closingParenIndex > filenameWithoutExtension.Length + 1) + { + var indexStr = entryNameWithoutExt.Substring( + filenameWithoutExtension.Length + 1, + closingParenIndex - filenameWithoutExtension.Length - 1); + + if (int.TryParse(indexStr, out int index)) + { + highestIndex = Math.Max(highestIndex, index); + } + } + } + } + + // Create a new name with the next available index + return $"{filenameWithoutExtension}({highestIndex + 1}){extension}"; + } } \ No newline at end of file diff --git a/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs b/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs index 45e82a262..8c4957b7d 100644 --- a/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs +++ b/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs @@ -96,27 +96,43 @@ public static class ContentTypeExtensions if (s.Length % 4 != 0) return false; - // Check valid Base64 characters - for (var i = 0; i < s.Length; i++) + // Check padding position and count + var paddingIndex = s.IndexOf('='); + if (paddingIndex > 0) + { + // Padding must be at the end + if (paddingIndex < s.Length - 2) + return false; + + // All characters after first '=' must also be '=' + if (s.Substring(paddingIndex).Any(c => c != '=')) + return false; + } + + // Check for valid Base64 characters + for (int i = 0; i < (paddingIndex > 0 ? paddingIndex : s.Length); i++) { var c = s[i]; - - var isValid = + var isValid = c is >= 'A' and <= 'Z' || c is >= 'a' and <= 'z' || c is >= '0' and <= '9' || - c == '+' || c == '/' || c == '='; + c == '+' || c == '/'; if (!isValid) return false; } - // Try actual decoding and roundtrip + // Additional check for short strings that are just lowercase+numbers + // This catches "content2" and similar false positives + if (s.Length <= 10 && s.All(c => char.IsLower(c) || char.IsDigit(c))) + return false; + + // Try actual decoding try { - var data = Convert.FromBase64String(s); - var reEncoded = Convert.ToBase64String(data); - return s == reEncoded; + _ = Convert.FromBase64String(s); + return true; } catch { From 792155c794b5db236c916f01618ca816140632be Mon Sep 17 00:00:00 2001 From: "lucas.hipolito" Date: Thu, 10 Jul 2025 10:24:32 +0200 Subject: [PATCH 04/14] Improving performance of existing entry check --- .../Activities/CreateZipArchive.cs | 56 +++++++++---------- 1 file changed, 28 insertions(+), 28 deletions(-) diff --git a/src/modules/Elsa.IO.Compression/Activities/CreateZipArchive.cs b/src/modules/Elsa.IO.Compression/Activities/CreateZipArchive.cs index 1e373731a..2b028549f 100644 --- a/src/modules/Elsa.IO.Compression/Activities/CreateZipArchive.cs +++ b/src/modules/Elsa.IO.Compression/Activities/CreateZipArchive.cs @@ -143,43 +143,39 @@ public class CreateZipArchive : CodeActivity private static string GetUniqueEntryName(ZipArchive zipArchive, string originalName) { - // If no duplicate exists, use the original name - if (!zipArchive.Entries.Any(entry => entry.Name.Equals(originalName, StringComparison.OrdinalIgnoreCase))) - { - return originalName; - } - - // Split the name into filename and extension - string filenameWithoutExtension = Path.GetFileNameWithoutExtension(originalName); - string extension = Path.GetExtension(originalName); - - // Find the highest index used for this filename pattern - int highestIndex = 0; + var filenameWithoutExtension = Path.GetFileNameWithoutExtension(originalName); + var extension = Path.GetExtension(originalName); + + var originalExists = false; + var highestIndex = 0; - // Check for the original name and any name with pattern "name(n).ext" foreach (var entry in zipArchive.Entries) { if (entry.Name.Equals(originalName, StringComparison.OrdinalIgnoreCase)) - continue; // Skip the exact match as we already know it exists - - string entryNameWithoutExt = Path.GetFileNameWithoutExtension(entry.Name); - string entryExt = Path.GetExtension(entry.Name); + { + originalExists = true; + } - if (!entryExt.Equals(extension, StringComparison.OrdinalIgnoreCase)) - continue; // Different extension + 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; - if (entryNameWithoutExt.StartsWith(filenameWithoutExtension, StringComparison.OrdinalIgnoreCase) && - entryNameWithoutExt.Length > filenameWithoutExtension.Length && - entryNameWithoutExt[filenameWithoutExtension.Length] == '(') + // Check if this entry follows our naming pattern + if (entryNameWithoutExtension.StartsWith(filenameWithoutExtension, StringComparison.OrdinalIgnoreCase) && + entryNameWithoutExtension.Length > filenameWithoutExtension.Length && + entryNameWithoutExtension[filenameWithoutExtension.Length] == '(') { // Extract the number between parentheses - var closingParenIndex = entryNameWithoutExt.LastIndexOf(')'); + var closingParenIndex = entryNameWithoutExtension.LastIndexOf(')'); if (closingParenIndex > filenameWithoutExtension.Length + 1) { - var indexStr = entryNameWithoutExt.Substring( - filenameWithoutExtension.Length + 1, + var indexStr = entryNameWithoutExtension.Substring( + filenameWithoutExtension.Length + 1, closingParenIndex - filenameWithoutExtension.Length - 1); - + if (int.TryParse(indexStr, out int index)) { highestIndex = Math.Max(highestIndex, index); @@ -187,8 +183,12 @@ public class CreateZipArchive : CodeActivity } } } - - // Create a new name with the next available index + + if (!originalExists) + { + return originalName; + } + return $"{filenameWithoutExtension}({highestIndex + 1}){extension}"; } } \ No newline at end of file From f6e1ed093a117da0d06da00a14674e0768ce1183 Mon Sep 17 00:00:00 2001 From: "lucas.hipolito" Date: Thu, 10 Jul 2025 10:44:35 +0200 Subject: [PATCH 05/14] small QoL improvements: more edge cases and constraints for base64 --- .../Activities/CreateZipArchive.cs | 3 +-- .../Elsa.IO/Extensions/ContentTypeExtensions.cs | 15 ++++++++------- 2 files changed, 9 insertions(+), 9 deletions(-) diff --git a/src/modules/Elsa.IO.Compression/Activities/CreateZipArchive.cs b/src/modules/Elsa.IO.Compression/Activities/CreateZipArchive.cs index 2b028549f..4c2b81ad6 100644 --- a/src/modules/Elsa.IO.Compression/Activities/CreateZipArchive.cs +++ b/src/modules/Elsa.IO.Compression/Activities/CreateZipArchive.cs @@ -125,8 +125,7 @@ public class CreateZipArchive : CodeActivity 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); diff --git a/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs b/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs index 8c4957b7d..9eac5a42a 100644 --- a/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs +++ b/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs @@ -98,19 +98,20 @@ public static class ContentTypeExtensions // Check padding position and count var paddingIndex = s.IndexOf('='); - if (paddingIndex > 0) + + switch (paddingIndex) { + // Padding cannot be at index 0 + case <= 0: // Padding must be at the end - if (paddingIndex < s.Length - 2) - return false; - + case > 0 when paddingIndex < s.Length - 2: // All characters after first '=' must also be '=' - if (s.Substring(paddingIndex).Any(c => c != '=')) + case > 0 when s[paddingIndex..].Any(c => c != '='): return false; } // Check for valid Base64 characters - for (int i = 0; i < (paddingIndex > 0 ? paddingIndex : s.Length); i++) + for (var i = 0; i < paddingIndex; i++) { var c = s[i]; var isValid = @@ -124,7 +125,7 @@ public static class ContentTypeExtensions } // Additional check for short strings that are just lowercase+numbers - // This catches "content2" and similar false positives + // This catches "whatever" and similar false positives if (s.Length <= 10 && s.All(c => char.IsLower(c) || char.IsDigit(c))) return false; From 9346595d80d6349422ebae9e129ac0204d70b941 Mon Sep 17 00:00:00 2001 From: "lucas.hipolito" Date: Thu, 10 Jul 2025 10:54:05 +0200 Subject: [PATCH 06/14] improving entry naming logic for better performance and readability --- .../Activities/CreateZipArchive.cs | 55 ++++++++++++------- 1 file changed, 35 insertions(+), 20 deletions(-) diff --git a/src/modules/Elsa.IO.Compression/Activities/CreateZipArchive.cs b/src/modules/Elsa.IO.Compression/Activities/CreateZipArchive.cs index 4c2b81ad6..e322d7594 100644 --- a/src/modules/Elsa.IO.Compression/Activities/CreateZipArchive.cs +++ b/src/modules/Elsa.IO.Compression/Activities/CreateZipArchive.cs @@ -126,6 +126,7 @@ public class CreateZipArchive : CodeActivity 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); @@ -150,11 +151,13 @@ public class CreateZipArchive : CodeActivity foreach (var entry in zipArchive.Entries) { - if (entry.Name.Equals(originalName, StringComparison.OrdinalIgnoreCase)) + if (!entry.Name.Equals(originalName, StringComparison.OrdinalIgnoreCase)) { - originalExists = true; + continue; } + originalExists = true; + var entryNameWithoutExtension = Path.GetFileNameWithoutExtension(entry.Name); var entryExtension = Path.GetExtension(entry.Name); @@ -163,24 +166,7 @@ public class CreateZipArchive : CodeActivity continue; // Check if this entry follows our naming pattern - if (entryNameWithoutExtension.StartsWith(filenameWithoutExtension, StringComparison.OrdinalIgnoreCase) && - entryNameWithoutExtension.Length > filenameWithoutExtension.Length && - entryNameWithoutExtension[filenameWithoutExtension.Length] == '(') - { - // Extract the number between parentheses - var closingParenIndex = entryNameWithoutExtension.LastIndexOf(')'); - if (closingParenIndex > filenameWithoutExtension.Length + 1) - { - var indexStr = entryNameWithoutExtension.Substring( - filenameWithoutExtension.Length + 1, - closingParenIndex - filenameWithoutExtension.Length - 1); - - if (int.TryParse(indexStr, out int index)) - { - highestIndex = Math.Max(highestIndex, index); - } - } - } + highestIndex = HighestEntryNameIndex(entryNameWithoutExtension, filenameWithoutExtension, highestIndex); } if (!originalExists) @@ -190,4 +176,33 @@ public class CreateZipArchive : CodeActivity 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 From fe08794f096ff1b121b7d7947ae0eed1cb923a2a Mon Sep 17 00:00:00 2001 From: "lucas.hipolito" Date: Thu, 10 Jul 2025 11:00:22 +0200 Subject: [PATCH 07/14] fix on padding index validation for base64 --- src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs b/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs index 9eac5a42a..45bf158ef 100644 --- a/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs +++ b/src/modules/Elsa.IO/Extensions/ContentTypeExtensions.cs @@ -102,7 +102,7 @@ public static class ContentTypeExtensions switch (paddingIndex) { // Padding cannot be at index 0 - case <= 0: + case 0: // Padding must be at the end case > 0 when paddingIndex < s.Length - 2: // All characters after first '=' must also be '=' From f9b9c84cf03e25eed761cb649cdcf092e73a83b6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Lucas=20Hip=C3=B3lito?= Date: Tue, 15 Jul 2025 13:59:39 +0200 Subject: [PATCH 08/14] Fixing Stream variable saved in instance issue (#6785) --- .../WorkflowInstanceStorageDriver.cs | 17 ++++++++++++++--- 1 file changed, 14 insertions(+), 3 deletions(-) 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; } From b535ed4b8c116ddd1bcd9c27e7b789729a8bfa0d Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 15 Jul 2025 14:09:48 +0200 Subject: [PATCH 09/14] Update Jint to 4.3.0 (#6794) --- Directory.Packages.props | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Directory.Packages.props b/Directory.Packages.props index a0b5e6b20..9a0667e87 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -59,7 +59,7 @@ - + From be8e61f1b4dacfb2197e9c6755088223ac8c2da0 Mon Sep 17 00:00:00 2001 From: "lucas.hipolito" Date: Tue, 15 Jul 2025 16:30:41 +0200 Subject: [PATCH 10/14] Implemented custom js functions: streamToBytes and streamToBase64 --- .../Handlers/ConfigureEngineWithCommonFunctions.cs | 11 ++++++++++- .../Providers/CommonFunctionsDefinitionProvider.cs | 12 +++++++++++- 2 files changed, 21 insertions(+), 2 deletions(-) diff --git a/src/modules/Elsa.JavaScript/Handlers/ConfigureEngineWithCommonFunctions.cs b/src/modules/Elsa.JavaScript/Handlers/ConfigureEngineWithCommonFunctions.cs index 0af980344..cd9ae1774 100644 --- a/src/modules/Elsa.JavaScript/Handlers/ConfigureEngineWithCommonFunctions.cs +++ b/src/modules/Elsa.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.JavaScript/Providers/CommonFunctionsDefinitionProvider.cs b/src/modules/Elsa.JavaScript/Providers/CommonFunctionsDefinitionProvider.cs index 76598428d..4922924cd 100644 --- a/src/modules/Elsa.JavaScript/Providers/CommonFunctionsDefinitionProvider.cs +++ b/src/modules/Elsa.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 From adbea90ccd7b4cd1fb2d922802699018e5732bf7 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Wed, 16 Jul 2025 18:05:32 +0200 Subject: [PATCH 11/14] Introduce activity execution record capturing and serialization improvements (#6800) * Introduce activity execution record capturing and serialization improvements - Added middleware for capturing activity execution records during workflow execution. - Introduced async mapping in `DefaultActivityExecutionMapper` with additional serialization support. - Enhanced `ActivityExecutionRecord` with new serialized properties for efficient storage. - Updated extensions to include `UseActivityExecutionLogCapturing`. - Simplified logging persistence by leveraging pre-serialized values in `ActivityExecutionLogStore`. * Update src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Update src/modules/Elsa.Workflows.Runtime/Middleware/Activities/CaptureActivityExecutionRecordMiddleware.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Apply suggestion from @Copilot Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --------- Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --- .../Runtime/ActivityExecutionLogStore.cs | 23 ++++------ .../DefaultActivityInvokerMiddleware.cs | 8 ---- .../Contracts/IActivityExecutonMapper.cs | 6 --- .../Entities/ActivityExecutionRecord.cs | 30 +++++++++---- ...ctivityExecutionContextEventExtensions.cs} | 2 +- ...ctivityExecutionContextRecordExtensions.cs | 23 ++++++++++ .../PipelineWorkflowsFeatureExtensions.cs | 1 + ...ivityExecutionPipelineBuilderExtensions.cs | 8 ++++ ...aptureActivityExecutionRecordMiddleware.cs | 13 ++++++ .../DefaultActivityExecutionMapper.cs | 43 ++++++++++++++----- .../Services/StoreActivityExecutionLogSink.cs | 3 +- 11 files changed, 110 insertions(+), 50 deletions(-) rename src/modules/Elsa.Workflows.Runtime/Extensions/{ActivityExecutionContextExtensions.cs => ActivityExecutionContextEventExtensions.cs} (96%) create mode 100644 src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Middleware/Activities/CaptureActivityExecutionRecordMiddleware.cs diff --git a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs index caeb89316..a71d19e6a 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs @@ -89,20 +89,13 @@ public class EFCoreActivityExecutionStore( private async ValueTask OnSaveAsync(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity, CancellationToken cancellationToken) { - entity = entity.SanitizeLogMessage(); - var compressionAlgorithm = options.Value.CompressionAlgorithm ?? nameof(None); - var serializedActivityState = entity.ActivityState?.Count > 0 ? safeSerializer.Serialize(entity.ActivityState) : null; - var compressedSerializedActivityState = serializedActivityState != null ? await compressionCodecResolver.Resolve(compressionAlgorithm).CompressAsync(serializedActivityState, cancellationToken) : null; - var serializedProperties = entity.Properties != null ? payloadSerializer.Serialize(entity.Properties) : null; - var serializedMetadata = entity.Metadata != null ? payloadSerializer.Serialize(entity.Metadata) : null; - - dbContext.Entry(entity).Property("SerializedActivityState").CurrentValue = compressedSerializedActivityState; - dbContext.Entry(entity).Property("SerializedActivityStateCompressionAlgorithm").CurrentValue = compressionAlgorithm; - dbContext.Entry(entity).Property("SerializedOutputs").CurrentValue = entity.Outputs?.Any() == true ? safeSerializer.Serialize(entity.Outputs) : null; - dbContext.Entry(entity).Property("SerializedProperties").CurrentValue = serializedProperties; - dbContext.Entry(entity).Property("SerializedMetadata").CurrentValue = serializedMetadata; - dbContext.Entry(entity).Property("SerializedException").CurrentValue = entity.Exception != null ? payloadSerializer.Serialize(entity.Exception) : null; - dbContext.Entry(entity).Property("SerializedPayload").CurrentValue = entity.Payload?.Any() == true ? payloadSerializer.Serialize(entity.Payload) : null; + dbContext.Entry(entity).Property("SerializedActivityState").CurrentValue = entity.SerializedActivityState; + dbContext.Entry(entity).Property("SerializedActivityStateCompressionAlgorithm").CurrentValue = entity.SerializedActivityStateCompressionAlgorithm ?? nameof(None); + dbContext.Entry(entity).Property("SerializedOutputs").CurrentValue = entity.SerializedOutputs; + dbContext.Entry(entity).Property("SerializedProperties").CurrentValue = entity.SerializedProperties; + dbContext.Entry(entity).Property("SerializedMetadata").CurrentValue = entity.SerializedMetadata; + dbContext.Entry(entity).Property("SerializedException").CurrentValue = entity.SerializedException; + dbContext.Entry(entity).Property("SerializedPayload").CurrentValue = entity.SerializedPayload; } [RequiresUnreferencedCode("Calls Elsa.EntityFrameworkCore.Modules.Runtime.EFCoreActivityExecutionStore.DeserializeActivityState(RuntimeElsaDbContext, ActivityExecutionRecord, CancellationToken)")] @@ -129,7 +122,7 @@ public class EFCoreActivityExecutionStore( var compressionAlgorithm = (string?)dbContext.Entry(entity).Property("SerializedActivityStateCompressionAlgorithm").CurrentValue ?? nameof(None); var compressionStrategy = compressionCodecResolver.Resolve(compressionAlgorithm); json = await compressionStrategy.DecompressAsync(json, cancellationToken); - var dictionary = JsonSerializer.Deserialize>(json); + var dictionary = safeSerializer.Deserialize?>(json); return dictionary?.ToDictionary(x => x.Key, x => x.Value); } diff --git a/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs b/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs index a4eb9ba8e..6731bb8a6 100644 --- a/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs +++ b/src/modules/Elsa.Workflows.Core/Middleware/Activities/DefaultActivityInvokerMiddleware.cs @@ -82,14 +82,6 @@ public class DefaultActivityInvokerMiddleware(ActivityMiddlewareDelegate next, I // Invoke next middleware. await next(context); - // // If the activity created any bookmarks, copy them into the workflow execution context. - // if (context.Bookmarks.Any()) - // { - // // Store bookmarks. - // workflowExecutionContext.Bookmarks.AddRange(context.Bookmarks); - // logger.LogDebug("Added {BookmarkCount} bookmarks to the workflow execution context", context.Bookmarks.Count); - // } - // Conditionally commit the workflow state. if (ShouldCommit(context, ActivityLifetimeEvent.ActivityExecuted)) await context.WorkflowExecutionContext.CommitAsync(); 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..fac785904 100644 --- a/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs +++ b/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs @@ -1,3 +1,4 @@ +using System.ComponentModel.DataAnnotations.Schema; using Elsa.Common; using Elsa.Common.Entities; using Elsa.Workflows.State; @@ -7,18 +8,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 +39,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 +59,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 +80,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 +90,20 @@ 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] public string? SerializedActivityState { get; set; } + [NotMapped] public string? SerializedOutputs { get; set; } + [NotMapped] public string? SerializedProperties { get; set; } + [NotMapped] public string? SerializedPayload { get; set; } + [NotMapped] public string? SerializedMetadata { get; set; } + [NotMapped] public string? SerializedException { get; set; } + [NotMapped] public string? SerializedActivityStateCompressionAlgorithm { 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..2dfb6f6ef --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs @@ -0,0 +1,23 @@ +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; + } +} \ 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/Services/DefaultActivityExecutionMapper.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs index e862e6be8..ea3de58c2 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs @@ -1,13 +1,24 @@ +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 +27,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 +50,26 @@ 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.SerializedActivityState = compressedSerializedActivityState; + record.SerializedActivityStateCompressionAlgorithm = compressionAlgorithm; + record.SerializedOutputs = record.Outputs?.Any() == true ? safeSerializer.Serialize(record.Outputs) : null; + record.SerializedProperties = serializedProperties; + record.SerializedMetadata = serializedMetadata; + record.SerializedException = record.Exception != null ? payloadSerializer.Serialize(record.Exception) : null; + record.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..9fad8be88 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(async x => x.GetCapturedActivityExecutionRecord() ?? await mapper.MapAsync(x))); await activityExecutionStore.SaveManyAsync(records, cancellationToken); // Untaint activity execution contexts. From b0fd60330a5d6c0bf71ea975ac2731c51d387334 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Wed, 16 Jul 2025 20:28:57 +0200 Subject: [PATCH 12/14] Simplify activity execution record retrieval with async mapping extension method --- .../ActivityExecutionContextRecordExtensions.cs | 9 +++++++++ .../Services/StoreActivityExecutionLogSink.cs | 2 +- 2 files changed, 10 insertions(+), 1 deletion(-) diff --git a/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs b/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs index 2dfb6f6ef..c2706cc1e 100644 --- a/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs +++ b/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs @@ -20,4 +20,13 @@ public static class ActivityExecutionContextRecordExtensions { 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/Services/StoreActivityExecutionLogSink.cs b/src/modules/Elsa.Workflows.Runtime/Services/StoreActivityExecutionLogSink.cs index 9fad8be88..72325955b 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/StoreActivityExecutionLogSink.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/StoreActivityExecutionLogSink.cs @@ -23,7 +23,7 @@ public class StoreActivityExecutionLogSink( if (activityExecutionContexts.Count == 0) return; - var records = await Task.WhenAll(activityExecutionContexts.Select(async x => x.GetCapturedActivityExecutionRecord() ?? await mapper.MapAsync(x))); + var records = await Task.WhenAll(activityExecutionContexts.Select(x => x.GetOrMapCapturedActivityExecutionRecordAsync())); await activityExecutionStore.SaveManyAsync(records, cancellationToken); // Untaint activity execution contexts. From 904284ce6b8846677f99ba6c15a95526daff2216 Mon Sep 17 00:00:00 2001 From: lukhipolito-nexxbiz Date: Fri, 18 Jul 2025 12:24:36 +0200 Subject: [PATCH 13/14] Fixed wrong naming in resulting entries when file is inputted from http (#6806) Co-authored-by: lucas.hipolito --- .../Services/Strategies/ZipEntryContentStrategy.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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; From ed14a1e577411bb8b7bf7ba8b4e2424a18e6a5a2 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 18 Jul 2025 14:14:19 +0200 Subject: [PATCH 14/14] Refactor activity execution record serialization with snapshots (#6807) - Introduced `ActivityExecutionRecordSnapshot` for encapsulated serialized data. - Updated `DefaultActivityExecutionMapper` to build serialized snapshots. - Adjusted `ActivityExecutionLogStore` to persist pre-serialized snapshots. - Streamlined package version management with `MicrosoftVersion` property. --- Directory.Packages.props | 67 ++++++++++--------- src/modules/Elsa.Common/Entities/Entity.cs | 2 +- .../Runtime/ActivityExecutionLogStore.cs | 29 ++++---- .../Entities/ActivityExecutionRecord.cs | 9 +-- .../Models/ActivityExecutionRecordSnapshot.cs | 25 +++++++ .../DefaultActivityExecutionMapper.cs | 48 ++++++++----- 6 files changed, 108 insertions(+), 72 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Runtime/Models/ActivityExecutionRecordSnapshot.cs diff --git a/Directory.Packages.props b/Directory.Packages.props index 9a0667e87..8b65fbfaa 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -5,6 +5,7 @@ 3.5.0-preview.1092 + 9.0.7 @@ -34,7 +35,7 @@ - + @@ -118,8 +119,8 @@ - - + + @@ -127,36 +128,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.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs index a71d19e6a..92b684ac9 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs @@ -6,17 +6,13 @@ using Elsa.Common.Codecs; using Elsa.Common.Entities; using Elsa.Extensions; using Elsa.Workflows; -using Elsa.Workflows.Management.Options; using Elsa.Workflows.Runtime; using Elsa.Workflows.Runtime.Entities; -using Elsa.Workflows.Runtime.Extensions; using Elsa.Workflows.Runtime.Filters; using Elsa.Workflows.Runtime.OrderDefinitions; using Elsa.Workflows.State; using JetBrains.Annotations; using Microsoft.EntityFrameworkCore; -using Microsoft.Extensions.Logging; -using Microsoft.Extensions.Options; using Open.Linq.AsyncExtensions; namespace Elsa.EntityFrameworkCore.Modules.Runtime; @@ -29,8 +25,7 @@ public class EFCoreActivityExecutionStore( EntityStore store, ISafeSerializer safeSerializer, IPayloadSerializer payloadSerializer, - ICompressionCodecResolver compressionCodecResolver, - IOptions options) : IActivityExecutionStore + ICompressionCodecResolver compressionCodecResolver) : IActivityExecutionStore { /// public async Task SaveAsync(ActivityExecutionRecord record, CancellationToken cancellationToken = default) => await store.SaveAsync(record, OnSaveAsync, cancellationToken); @@ -87,15 +82,21 @@ public class EFCoreActivityExecutionStore( return await store.DeleteWhereAsync(queryable => Filter(queryable, filter), cancellationToken); } - private async ValueTask OnSaveAsync(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity, CancellationToken cancellationToken) + private ValueTask OnSaveAsync(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity, CancellationToken cancellationToken) { - dbContext.Entry(entity).Property("SerializedActivityState").CurrentValue = entity.SerializedActivityState; - dbContext.Entry(entity).Property("SerializedActivityStateCompressionAlgorithm").CurrentValue = entity.SerializedActivityStateCompressionAlgorithm ?? nameof(None); - dbContext.Entry(entity).Property("SerializedOutputs").CurrentValue = entity.SerializedOutputs; - dbContext.Entry(entity).Property("SerializedProperties").CurrentValue = entity.SerializedProperties; - dbContext.Entry(entity).Property("SerializedMetadata").CurrentValue = entity.SerializedMetadata; - dbContext.Entry(entity).Property("SerializedException").CurrentValue = entity.SerializedException; - dbContext.Entry(entity).Property("SerializedPayload").CurrentValue = entity.SerializedPayload; + var snapshot = entity.SerializedSnapshot; + + if (snapshot is null) + return ValueTask.CompletedTask; + + dbContext.Entry(entity).Property("SerializedActivityState").CurrentValue = snapshot.SerializedActivityState; + dbContext.Entry(entity).Property("SerializedActivityStateCompressionAlgorithm").CurrentValue = snapshot.SerializedActivityStateCompressionAlgorithm; + dbContext.Entry(entity).Property("SerializedOutputs").CurrentValue = snapshot.SerializedOutputs; + dbContext.Entry(entity).Property("SerializedProperties").CurrentValue = snapshot.SerializedProperties; + dbContext.Entry(entity).Property("SerializedMetadata").CurrentValue = snapshot.SerializedMetadata; + dbContext.Entry(entity).Property("SerializedException").CurrentValue = snapshot.SerializedException; + dbContext.Entry(entity).Property("SerializedPayload").CurrentValue = snapshot.SerializedPayload; + return ValueTask.CompletedTask; } [RequiresUnreferencedCode("Calls Elsa.EntityFrameworkCore.Modules.Runtime.EFCoreActivityExecutionStore.DeserializeActivityState(RuntimeElsaDbContext, ActivityExecutionRecord, CancellationToken)")] diff --git a/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs b/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs index fac785904..2b6bd1e74 100644 --- a/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs +++ b/src/modules/Elsa.Workflows.Runtime/Entities/ActivityExecutionRecord.cs @@ -1,4 +1,5 @@ using System.ComponentModel.DataAnnotations.Schema; +using System.Text.Json.Serialization; using Elsa.Common; using Elsa.Common.Entities; using Elsa.Workflows.State; @@ -99,11 +100,5 @@ public partial class ActivityExecutionRecord : Entity, ILogRecord public partial class ActivityExecutionRecord { - [NotMapped] public string? SerializedActivityState { get; set; } - [NotMapped] public string? SerializedOutputs { get; set; } - [NotMapped] public string? SerializedProperties { get; set; } - [NotMapped] public string? SerializedPayload { get; set; } - [NotMapped] public string? SerializedMetadata { get; set; } - [NotMapped] public string? SerializedException { get; set; } - [NotMapped] public string? SerializedActivityStateCompressionAlgorithm { get; set; } + [NotMapped] [JsonIgnore] public ActivityExecutionRecordSnapshot? SerializedSnapshot { get; set; } } \ 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 ea3de58c2..223cd55f3 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs @@ -11,12 +11,11 @@ namespace Elsa.Workflows.Runtime; /// public class DefaultActivityExecutionMapper( - ISafeSerializer safeSerializer, + ISafeSerializer safeSerializer, IPayloadSerializer payloadSerializer, ICompressionCodecResolver compressionCodecResolver, IOptions options) : IActivityExecutionMapper { - /// public async Task MapAsync(ActivityExecutionContext source) { @@ -50,22 +49,37 @@ public class DefaultActivityExecutionMapper( 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.SerializedActivityState = compressedSerializedActivityState; - record.SerializedActivityStateCompressionAlgorithm = compressionAlgorithm; - record.SerializedOutputs = record.Outputs?.Any() == true ? safeSerializer.Serialize(record.Outputs) : null; - record.SerializedProperties = serializedProperties; - record.SerializedMetadata = serializedMetadata; - record.SerializedException = record.Exception != null ? payloadSerializer.Serialize(record.Exception) : null; - record.SerializedPayload = record.Payload?.Any() == true ? payloadSerializer.Serialize(record.Payload) : null; + 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; }