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