From adbea90ccd7b4cd1fb2d922802699018e5732bf7 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Wed, 16 Jul 2025 18:05:32 +0200 Subject: [PATCH] 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.