From c8ed08cc850decae8f4ea736af3a5ee17c92cfa1 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 18 Apr 2025 16:21:25 +0200 Subject: [PATCH] Remove obsolete properties and refactor activity evaluation (#6603) * Remove obsolete properties and refactor activity evaluation Refactored activity input and log persistence property evaluation using improved notification handlers. Removed redundant `ActivityState` property and associated serialization logic, ensuring payloads are only serialized when necessary. All changes streamline workflow processing and enhance maintainability. * Refactor mediator call to inline cancellation token. Replaced the separate variable for the cancellation token with an inline reference for clarity and reduced redundancy. This simplifies the code without altering functionality. --- .../Stores/DapperWorkflowExecutionLogStore.cs | 20 +++++++++++++------ .../Runtime/WorkflowExecutionLogStore.cs | 14 +++++++------ .../Elsa.Scheduling/Activities/Cron.cs | 10 +++++----- .../Models/ExecutionLogRecord.cs | 2 +- .../Abstractions/Activity.cs | 8 ++++---- .../ScheduledChildCallbackBehavior.cs | 19 +++++++++--------- .../ActivityExecutionContext.Cancel.cs | 2 +- .../ActivityExecutionContext.Complete.cs | 2 +- .../ActivityExecutionContextExtensions.cs | 2 +- .../Handlers/EvaluateParentInputProperties.cs | 15 ++++++++++++++ .../Activities/ExecutionLogMiddleware.cs | 2 +- .../Notifications/InvokingActivityCallback.cs | 5 +++++ .../Entities/WorkflowExecutionLogRecord.cs | 1 + .../Features/WorkflowRuntimeFeature.cs | 1 + .../ActivityExecutionContextExtensions.cs | 5 +++++ .../Handlers/EvaluateParentInputProperties.cs | 17 ++++++++++++++++ .../DefaultActivityExecutionMapper.cs | 16 ++------------- .../WorkflowExecutionLogRecordExtractor.cs | 1 - 18 files changed, 92 insertions(+), 50 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Core/Handlers/EvaluateParentInputProperties.cs create mode 100644 src/modules/Elsa.Workflows.Core/Notifications/InvokingActivityCallback.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/LogPersistence/Handlers/EvaluateParentInputProperties.cs diff --git a/src/modules/Elsa.Dapper/Modules/Runtime/Stores/DapperWorkflowExecutionLogStore.cs b/src/modules/Elsa.Dapper/Modules/Runtime/Stores/DapperWorkflowExecutionLogStore.cs index 5e5dde8f3..32d3db0fb 100644 --- a/src/modules/Elsa.Dapper/Modules/Runtime/Stores/DapperWorkflowExecutionLogStore.cs +++ b/src/modules/Elsa.Dapper/Modules/Runtime/Stores/DapperWorkflowExecutionLogStore.cs @@ -104,12 +104,12 @@ internal class DapperWorkflowExecutionLogStore(Store Map(Page source) { - return new Page(source.Items.Select(Map).ToList(), source.TotalCount); + return new(source.Items.Select(Map).ToList(), source.TotalCount); } private WorkflowExecutionLogRecordRecord Map(WorkflowExecutionLogRecord source) { - return new WorkflowExecutionLogRecordRecord + return new() { Id = source.Id, WorkflowDefinitionId = source.WorkflowDefinitionId, @@ -128,15 +128,24 @@ internal class DapperWorkflowExecutionLogStore(Store false, + IDictionary dictionary => dictionary.Count > 0, + _ => true + }; + } + private WorkflowExecutionLogRecord Map(WorkflowExecutionLogRecordRecord source) { - return new WorkflowExecutionLogRecord + return new() { Id = source.Id, WorkflowDefinitionId = source.WorkflowDefinitionId, @@ -155,7 +164,6 @@ internal class DapperWorkflowExecutionLogStore(Store>(source.SerializedActivityState) : null, Payload = source.SerializedPayload != null ? payloadSerializer.Deserialize(source.SerializedPayload) : null, TenantId = source.TenantId }; diff --git a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/WorkflowExecutionLogStore.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/WorkflowExecutionLogStore.cs index 16bd0a05c..03b9088cb 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/WorkflowExecutionLogStore.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/WorkflowExecutionLogStore.cs @@ -74,8 +74,7 @@ public class EFCoreWorkflowExecutionLogStore(EntityStore LoadPayload(RuntimeElsaDbContext dbContext, WorkflowExecutionLogRecord entity) @@ -93,10 +91,14 @@ public class EFCoreWorkflowExecutionLogStore(EntityStore(json) : null); } - private ValueTask?> LoadActivityState(RuntimeElsaDbContext dbContext, WorkflowExecutionLogRecord entity) + private bool ShouldSerializePayload(WorkflowExecutionLogRecord source) { - var json = dbContext.Entry(entity).Property("SerializedActivityState").CurrentValue; - return new(!string.IsNullOrEmpty(json) ? JsonSerializer.Deserialize>(json) : null); + return source.Payload switch + { + null => false, + IDictionary dictionary => dictionary.Count > 0, + _ => true + }; } private static IQueryable Filter(IQueryable queryable, WorkflowExecutionLogRecordFilter filter) => filter.Apply(queryable); diff --git a/src/modules/Elsa.Scheduling/Activities/Cron.cs b/src/modules/Elsa.Scheduling/Activities/Cron.cs index 1f5976095..6f0f84012 100644 --- a/src/modules/Elsa.Scheduling/Activities/Cron.cs +++ b/src/modules/Elsa.Scheduling/Activities/Cron.cs @@ -14,17 +14,17 @@ namespace Elsa.Scheduling.Activities; public class Cron : EventGenerator { /// - public Cron([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line) + public Cron([CallerFilePath] string? source = null, [CallerLineNumber] int? line = null) : base(source, line) { } /// - public Cron(string cronExpression, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(new Input(cronExpression), source, line) + public Cron(string cronExpression, [CallerFilePath] string? source = null, [CallerLineNumber] int? line = null) : this(new Input(cronExpression), source, line) { } /// - public Cron(Input cronExpression, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : this(source, line) + public Cron(Input cronExpression, [CallerFilePath] string? source = null, [CallerLineNumber] int? line = null) : this(source, line) { CronExpression = cronExpression; } @@ -33,7 +33,7 @@ public class Cron : EventGenerator /// The interval at which the timer should execute. /// [Input(Description = "The CRON expression at which the timer should execute.")] - public Input CronExpression { get; set; } = default!; + public Input CronExpression { get; set; } = null!; /// protected override async ValueTask ExecuteAsync(ActivityExecutionContext context) @@ -62,5 +62,5 @@ public class Cron : EventGenerator /// /// Creates a new activity set to trigger at the specified cron expression. /// - public static Cron FromCronExpression(string value, [CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) => new(value, source, line); + public static Cron FromCronExpression(string value, [CallerFilePath] string? source = null, [CallerLineNumber] int? line = null) => new(value, source, line); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Api/Models/ExecutionLogRecord.cs b/src/modules/Elsa.Workflows.Api/Models/ExecutionLogRecord.cs index 492f6c588..85fe32110 100644 --- a/src/modules/Elsa.Workflows.Api/Models/ExecutionLogRecord.cs +++ b/src/modules/Elsa.Workflows.Api/Models/ExecutionLogRecord.cs @@ -14,5 +14,5 @@ internal record ExecutionLogRecord( string? EventName, string? Message, string? Source, - IDictionary? ActivityState, + [property: Obsolete] IDictionary? ActivityState, object? Payload); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Abstractions/Activity.cs b/src/modules/Elsa.Workflows.Core/Abstractions/Activity.cs index 3be660125..0ef11e6d8 100644 --- a/src/modules/Elsa.Workflows.Core/Abstractions/Activity.cs +++ b/src/modules/Elsa.Workflows.Core/Abstractions/Activity.cs @@ -21,7 +21,7 @@ public abstract class Activity : IActivity, ISignalHandler /// /// Constructor. /// - protected Activity(string? source = default, int? line = default) + protected Activity(string? source = null, int? line = null) { this.SetSource(source, line); Type = ActivityTypeNameHelper.GenerateTypeName(GetType()); @@ -30,17 +30,17 @@ public abstract class Activity : IActivity, ISignalHandler } /// - protected Activity(string activityType, int version = 1, string? source = default, int? line = default) : this(source, line) + protected Activity(string activityType, int version = 1, string? source = null, int? line = null) : this(source, line) { Type = activityType; Version = version; } /// - public string Id { get; set; } = default!; + public string Id { get; set; } = null!; /// - public string NodeId { get; set; } = default!; + public string NodeId { get; set; } = null!; /// public string? Name { get; set; } diff --git a/src/modules/Elsa.Workflows.Core/Behaviors/ScheduledChildCallbackBehavior.cs b/src/modules/Elsa.Workflows.Core/Behaviors/ScheduledChildCallbackBehavior.cs index 1dd1b6d94..709f58cfd 100644 --- a/src/modules/Elsa.Workflows.Core/Behaviors/ScheduledChildCallbackBehavior.cs +++ b/src/modules/Elsa.Workflows.Core/Behaviors/ScheduledChildCallbackBehavior.cs @@ -1,11 +1,14 @@ -using Elsa.Extensions; +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Notifications; using Elsa.Workflows.Signals; +using JetBrains.Annotations; namespace Elsa.Workflows.Behaviors; /// /// Implements a behavior that invokes "child completed" callbacks on parent activities. /// +[UsedImplicitly] public class ScheduledChildCallbackBehavior : Behavior { /// @@ -21,18 +24,16 @@ public class ScheduledChildCallbackBehavior : Behavior var childActivityNode = childActivityExecutionContext.ActivityNode; var callbackEntry = activityExecutionContext.WorkflowExecutionContext.PopCompletionCallback(activityExecutionContext, childActivityNode); - if (callbackEntry == null) - return; - - // Before invoking the parent activity, make sure its properties are evaluated. - if (!activityExecutionContext.GetHasEvaluatedProperties()) - await activityExecutionContext.EvaluateInputPropertiesAsync(); - - if (callbackEntry.CompletionCallback != null) + if (callbackEntry?.CompletionCallback != null) { var completedContext = new ActivityCompletedContext(activityExecutionContext, childActivityExecutionContext, signal.Result); var tag = callbackEntry.Tag; completedContext.TargetContext.Tag = tag; + + var mediator = activityExecutionContext.GetRequiredService(); + var invokingActivityCallbackNotification = new InvokingActivityCallback(activityExecutionContext, childActivityExecutionContext); + await mediator.SendAsync(invokingActivityCallbackNotification, context.CancellationToken); + await callbackEntry.CompletionCallback(completedContext); } } diff --git a/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.Cancel.cs b/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.Cancel.cs index 84ddc20c1..1ad0adbe2 100644 --- a/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.Cancel.cs +++ b/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.Cancel.cs @@ -23,7 +23,7 @@ public partial class ActivityExecutionContext ClearBookmarks(); ClearCompletionCallbacks(); WorkflowExecutionContext.Bookmarks.RemoveWhere(x => x.ActivityNodeId == NodeId); - AddExecutionLogEntry("Canceled", payload: JournalData); + AddExecutionLogEntry("Canceled"); await this.SendSignalAsync(new CancelSignal()); await CancelChildActivitiesAsync(); diff --git a/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.Complete.cs b/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.Complete.cs index 1a278d1db..5379093ff 100644 --- a/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.Complete.cs +++ b/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.Complete.cs @@ -41,7 +41,7 @@ public partial class ActivityExecutionContext JournalData["Outcomes"] = outcomes.Names; // Add an execution log entry. - AddExecutionLogEntry("Completed", payload: JournalData); + AddExecutionLogEntry("Completed"); // Send a signal. await this.SendSignalAsync(new ActivityCompleted(result)); diff --git a/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs b/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs index f229bee04..9be4d693d 100644 --- a/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs +++ b/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs @@ -294,7 +294,7 @@ public static partial class ActivityExecutionContextExtensions context.WorkflowExecutionContext.Bookmarks.RemoveWhere(x => x.ActivityNodeId == context.NodeId); // Add an execution log entry. - context.AddExecutionLogEntry("Canceled", payload: context.JournalData); + context.AddExecutionLogEntry("Canceled"); await context.SendSignalAsync(new CancelSignal()); await publisher.SendAsync(new ActivityCancelled(context)); diff --git a/src/modules/Elsa.Workflows.Core/Handlers/EvaluateParentInputProperties.cs b/src/modules/Elsa.Workflows.Core/Handlers/EvaluateParentInputProperties.cs new file mode 100644 index 000000000..b7714e447 --- /dev/null +++ b/src/modules/Elsa.Workflows.Core/Handlers/EvaluateParentInputProperties.cs @@ -0,0 +1,15 @@ +using Elsa.Extensions; +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Notifications; + +namespace Elsa.Workflows.Handlers; + +public class EvaluateParentInputProperties : INotificationHandler +{ + public async Task HandleAsync(InvokingActivityCallback notification, CancellationToken cancellationToken) + { + // Before invoking the parent activity, make sure its properties are evaluated. + if (!notification.Parent.GetHasEvaluatedProperties()) + await notification.Parent.EvaluateInputPropertiesAsync(); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Middleware/Activities/ExecutionLogMiddleware.cs b/src/modules/Elsa.Workflows.Core/Middleware/Activities/ExecutionLogMiddleware.cs index 041148ef1..50f63e6f9 100644 --- a/src/modules/Elsa.Workflows.Core/Middleware/Activities/ExecutionLogMiddleware.cs +++ b/src/modules/Elsa.Workflows.Core/Middleware/Activities/ExecutionLogMiddleware.cs @@ -33,7 +33,7 @@ public class ExecutionLogMiddleware(ActivityMiddlewareDelegate next) : IActivity if (context.Status == ActivityStatus.Running) { if (IsActivityBookmarked(context)) - context.AddExecutionLogEntry("Suspended", payload: context.JournalData); + context.AddExecutionLogEntry("Suspended"); } } catch (Exception exception) diff --git a/src/modules/Elsa.Workflows.Core/Notifications/InvokingActivityCallback.cs b/src/modules/Elsa.Workflows.Core/Notifications/InvokingActivityCallback.cs new file mode 100644 index 000000000..9d2d133cc --- /dev/null +++ b/src/modules/Elsa.Workflows.Core/Notifications/InvokingActivityCallback.cs @@ -0,0 +1,5 @@ +using Elsa.Mediator.Contracts; + +namespace Elsa.Workflows.Notifications; + +public record InvokingActivityCallback(ActivityExecutionContext Parent, ActivityExecutionContext Child) : INotification; \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Entities/WorkflowExecutionLogRecord.cs b/src/modules/Elsa.Workflows.Runtime/Entities/WorkflowExecutionLogRecord.cs index 54c412db1..cf256dd11 100644 --- a/src/modules/Elsa.Workflows.Runtime/Entities/WorkflowExecutionLogRecord.cs +++ b/src/modules/Elsa.Workflows.Runtime/Entities/WorkflowExecutionLogRecord.cs @@ -91,6 +91,7 @@ public class WorkflowExecutionLogRecord : Entity, ILogRecord /// /// The state of the activity at the time of the log entry. /// + [Obsolete("Look at the ActivityExecutionRecord.ActivityState property instead.")] public IDictionary? ActivityState { get; set; } /// diff --git a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs index caa7aba26..7566d57a0 100644 --- a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs @@ -344,6 +344,7 @@ public class WorkflowRuntimeFeature(IModule module) : FeatureBase(module) .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() + .AddNotificationHandler() // Workflow activation strategies. .AddScoped() diff --git a/src/modules/Elsa.Workflows.Runtime/LogPersistence/Extensions/ActivityExecutionContextExtensions.cs b/src/modules/Elsa.Workflows.Runtime/LogPersistence/Extensions/ActivityExecutionContextExtensions.cs index 1d6ce8577..fe9edd3a3 100644 --- a/src/modules/Elsa.Workflows.Runtime/LogPersistence/Extensions/ActivityExecutionContextExtensions.cs +++ b/src/modules/Elsa.Workflows.Runtime/LogPersistence/Extensions/ActivityExecutionContextExtensions.cs @@ -6,6 +6,11 @@ public static class ActivityExecutionContextExtensions { private static object LogPersistenceMapKey { get; } = new(); + public static bool HasLogPersistenceModeMap(this ActivityExecutionContext context) + { + return context.TransientProperties.ContainsKey(LogPersistenceMapKey); + } + public static ActivityLogPersistenceModeMap GetLogPersistenceModeMap(this ActivityExecutionContext context) { return context.TransientProperties.GetValueOrDefault(LogPersistenceMapKey, () => new ActivityLogPersistenceModeMap())!; diff --git a/src/modules/Elsa.Workflows.Runtime/LogPersistence/Handlers/EvaluateParentInputProperties.cs b/src/modules/Elsa.Workflows.Runtime/LogPersistence/Handlers/EvaluateParentInputProperties.cs new file mode 100644 index 000000000..ff93d4b0e --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/LogPersistence/Handlers/EvaluateParentInputProperties.cs @@ -0,0 +1,17 @@ +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Notifications; + +namespace Elsa.Workflows.Runtime.Handlers; + +public class EvaluateParentLogPersistenceModes(IActivityPropertyLogPersistenceEvaluator persistenceEvaluator) : INotificationHandler +{ + public async Task HandleAsync(InvokingActivityCallback notification, CancellationToken cancellationToken) + { + // Before invoking the parent activity, make sure its persistence log properties are evaluated. + if (!notification.Parent.HasLogPersistenceModeMap()) + { + var persistenceLogMap = await persistenceEvaluator.EvaluateLogPersistenceModesAsync(notification.Parent); + notification.Parent.SetLogPersistenceModeMap(persistenceLogMap); + } + } +} \ 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 054da084f..20744b097 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs @@ -9,14 +9,13 @@ public class DefaultActivityExecutionMapper : IActivityExecutionMapper { public ActivityExecutionRecord Map(ActivityExecutionContext source) { - var payload = GetPayload(source); var outputs = source.GetOutputs(); var inputs = source.GetInputs(); var persistenceMap = source.GetLogPersistenceModeMap(); var persistableInputs = GetPersistableInputOutput(inputs, persistenceMap.Inputs); var persistableOutputs = GetPersistableInputOutput(outputs, persistenceMap.Outputs); var persistableProperties = GetPersistableDictionary(source.Properties!, persistenceMap.InternalState); - var persistablePayload = GetPersistableDictionary(payload!, persistenceMap.InternalState); + var persistableJournalData = GetPersistableDictionary(source.JournalData!, persistenceMap.InternalState); return new() { @@ -29,7 +28,7 @@ public class DefaultActivityExecutionMapper : IActivityExecutionMapper ActivityState = persistableInputs, Outputs = persistableOutputs, Properties = persistableProperties, - Payload = persistablePayload!, + Payload = persistableJournalData!, Exception = ExceptionState.FromException(source.Exception), ActivityTypeVersion = source.Activity.Version, StartedAt = source.StartedAt, @@ -62,15 +61,4 @@ public class DefaultActivityExecutionMapper : IActivityExecutionMapper { return mode == LogPersistenceMode.Include ? dictionary : null; } - - private static IDictionary GetPayload(ActivityExecutionContext source) - { - var outcomes = source.JournalData.TryGetValue("Outcomes", out var resultValue) ? resultValue as string[] : null; - var payload = new Dictionary(); - - if (outcomes != null) - payload.Add("Outcomes", outcomes); - - return payload; - } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionLogRecordExtractor.cs b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionLogRecordExtractor.cs index 081016122..35f9fc174 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionLogRecordExtractor.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowExecutionLogRecordExtractor.cs @@ -25,7 +25,6 @@ public class WorkflowExecutionLogRecordExtractor(IIdentityGenerator identityGene WorkflowInstanceId = context.Id, WorkflowVersion = context.Workflow.Version, Source = x.Source, - ActivityState = x.ActivityState, Payload = x.Payload, Timestamp = x.Timestamp, Sequence = x.Sequence