From 5a2b5396378387b07008e39e226e434d69f7cdb8 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 3 Jun 2022 14:22:55 +0200 Subject: [PATCH] Implement Started and Completed execution events --- .../Activities/Flowchart.cs | 2 +- .../Activities/Sequence.cs | 2 +- .../Behaviors/BreakBehavior.cs | 6 ++-- .../Behaviors/ExecutionLoggingBehavior.cs | 28 +++++++++++++++++++ .../ScheduledChildCallbackBehavior.cs | 4 +-- .../ActivityExecutionContextExtensions.cs | 15 ++++------ .../Elsa.Workflows.Core/Models/Activity.cs | 2 +- .../Components/ActivityInvokerMiddleware.cs | 7 ++--- .../Services/ISignalHandler.cs | 23 +++++++++++---- 9 files changed, 62 insertions(+), 27 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Core/Behaviors/ExecutionLoggingBehavior.cs diff --git a/src/modules/Elsa.Workflows.Core/Activities/Flowchart.cs b/src/modules/Elsa.Workflows.Core/Activities/Flowchart.cs index ee1c5f781..704b54eb3 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart.cs @@ -28,7 +28,7 @@ public class Flowchart : Container private async ValueTask OnDescendantCompletedAsync(ActivityCompleted signal, SignalContext context) { - await ScheduleChildrenAsync(context.ActivityExecutionContext, context.SourceActivityExecutionContext.Activity); + await ScheduleChildrenAsync(context.ReceiverActivityExecutionContext, context.SenderActivityExecutionContext.Activity); } private async Task ScheduleChildrenAsync(ActivityExecutionContext context, IActivity parent) diff --git a/src/modules/Elsa.Workflows.Core/Activities/Sequence.cs b/src/modules/Elsa.Workflows.Core/Activities/Sequence.cs index e9adb7690..d82c51d4c 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Sequence.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Sequence.cs @@ -54,6 +54,6 @@ public class Sequence : Container private void OnBreak(BreakSignal signal, SignalContext context) { // Clear any scheduled child completion callbacks, since we no longer want to schedule any sibling. - context.ActivityExecutionContext.ClearCompletionCallbacks(); + context.ReceiverActivityExecutionContext.ClearCompletionCallbacks(); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Behaviors/BreakBehavior.cs b/src/modules/Elsa.Workflows.Core/Behaviors/BreakBehavior.cs index c1631a62e..0efa4721c 100644 --- a/src/modules/Elsa.Workflows.Core/Behaviors/BreakBehavior.cs +++ b/src/modules/Elsa.Workflows.Core/Behaviors/BreakBehavior.cs @@ -20,10 +20,10 @@ public class BreakBehavior : Behavior context.StopPropagation(); // Remove child activity execution contexts. - var childActivityExecutionContexts = context.ActivityExecutionContext.GetChildren().ToList(); - context.ActivityExecutionContext.WorkflowExecutionContext.RemoveActivityExecutionContexts(childActivityExecutionContexts); + var childActivityExecutionContexts = context.ReceiverActivityExecutionContext.GetChildren().ToList(); + context.ReceiverActivityExecutionContext.WorkflowExecutionContext.RemoveActivityExecutionContexts(childActivityExecutionContexts); // Mark this activity as completed. - await context.ActivityExecutionContext.CompleteActivityAsync(); + await context.ReceiverActivityExecutionContext.CompleteActivityAsync(); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Behaviors/ExecutionLoggingBehavior.cs b/src/modules/Elsa.Workflows.Core/Behaviors/ExecutionLoggingBehavior.cs new file mode 100644 index 000000000..0f0948bdd --- /dev/null +++ b/src/modules/Elsa.Workflows.Core/Behaviors/ExecutionLoggingBehavior.cs @@ -0,0 +1,28 @@ +using System.Runtime.CompilerServices; +using Elsa.Workflows.Core.Models; +using Elsa.Workflows.Core.Services; +using Elsa.Workflows.Core.Signals; + +namespace Elsa.Workflows.Core.Behaviors; + +/// +/// Records Start and Completed entries of the owning activity. +/// +public class ExecutionLoggingBehavior : Behavior +{ + public ExecutionLoggingBehavior(IActivity owner) : base(owner) + { + OnSignalReceived(OnActivityCompleted); + } + + protected override void Execute(ActivityExecutionContext context) + { + context.AddExecutionLogEntry("Started"); + } + + private void OnActivityCompleted(ActivityCompleted signal, SignalContext context) + { + if (context.IsSelf) + context.SenderActivityExecutionContext.AddExecutionLogEntry("Completed"); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Behaviors/ScheduledChildCallbackBehavior.cs b/src/modules/Elsa.Workflows.Core/Behaviors/ScheduledChildCallbackBehavior.cs index 97c14e383..9e70d0e38 100644 --- a/src/modules/Elsa.Workflows.Core/Behaviors/ScheduledChildCallbackBehavior.cs +++ b/src/modules/Elsa.Workflows.Core/Behaviors/ScheduledChildCallbackBehavior.cs @@ -16,8 +16,8 @@ public class ScheduledChildCallbackBehavior : Behavior private async ValueTask OnActivityCompletedAsync(ActivityCompleted signal, SignalContext context) { - var activityExecutionContext = context.ActivityExecutionContext; - var childActivityExecutionContext = context.SourceActivityExecutionContext; + var activityExecutionContext = context.ReceiverActivityExecutionContext; + var childActivityExecutionContext = context.SenderActivityExecutionContext; var childActivity = childActivityExecutionContext.Activity; var callbackEntry = activityExecutionContext.WorkflowExecutionContext.PopCompletionCallback(activityExecutionContext, childActivity); diff --git a/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs b/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs index db87a6369..3473bdabc 100644 --- a/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs +++ b/src/modules/Elsa.Workflows.Core/Extensions/ActivityExecutionContextExtensions.cs @@ -13,9 +13,6 @@ public static class ActivityExecutionContextExtensions public static T GetInput(this ActivityExecutionContext context) => context.GetInput(typeof(T).Name); public static T GetInput(this ActivityExecutionContext context, string key) => (T)context.Input[key]; - public static WorkflowExecutionLogEntry AddExecutionLogEntry(this ActivityExecutionContext context, string eventName, string? message = default, object? payload = default) => - context.AddExecutionLogEntry(eventName, message, default, payload); - public static WorkflowExecutionLogEntry AddExecutionLogEntry(this ActivityExecutionContext context, string eventName, string? message = default, string? source = default, object? payload = default) { var activity = context.Activity; @@ -78,7 +75,7 @@ public static class ActivityExecutionContextExtensions return input; } - + public static async Task EvaluateAsync(this ActivityExecutionContext context, Input input) { var evaluator = context.GetRequiredService(); @@ -87,7 +84,7 @@ public static class ActivityExecutionContextExtensions locationReference.Set(context, value); return value; } - + /// /// Returns a flattened list of the current context's ancestors. /// @@ -107,7 +104,7 @@ public static class ActivityExecutionContextExtensions /// Returns a flattened list of the current context's immediate children. /// /// - public static IEnumerable GetChildren(this ActivityExecutionContext context) => + public static IEnumerable GetChildren(this ActivityExecutionContext context) => context.WorkflowExecutionContext.ActivityExecutionContexts.Where(x => x.ParentActivityExecutionContext == context); /// @@ -118,13 +115,13 @@ public static class ActivityExecutionContextExtensions // Detach child activity execution contexts. context.WorkflowExecutionContext.RemoveActivityExecutionContexts(context.GetChildren()); } - + /// /// Send a signal up the current branch. /// public static async ValueTask SignalAsync(this ActivityExecutionContext context, object signal) { - var ancestorContexts = context.GetAncestors(); + var ancestorContexts = new[] { context }.Concat(context.GetAncestors()); foreach (var ancestorContext in ancestorContexts) { @@ -132,7 +129,7 @@ public static class ActivityExecutionContextExtensions if (ancestorContext.Activity is not ISignalHandler handler) continue; - + await handler.HandleSignalAsync(signal, signalContext); if (signalContext.StopPropagationRequested) diff --git a/src/modules/Elsa.Workflows.Core/Models/Activity.cs b/src/modules/Elsa.Workflows.Core/Models/Activity.cs index ad325dcb6..f71af1a19 100644 --- a/src/modules/Elsa.Workflows.Core/Models/Activity.cs +++ b/src/modules/Elsa.Workflows.Core/Models/Activity.cs @@ -14,9 +14,9 @@ public abstract class Activity : IActivity, ISignalHandler protected Activity() { TypeName = ActivityTypeNameHelper.GenerateTypeName(GetType()); + Behaviors.Add(this); Behaviors.Add(this); Behaviors.Add(this); - Behaviors.Add(this); } protected Activity(string activityType) : this() diff --git a/src/modules/Elsa.Workflows.Core/Pipelines/ActivityExecution/Components/ActivityInvokerMiddleware.cs b/src/modules/Elsa.Workflows.Core/Pipelines/ActivityExecution/Components/ActivityInvokerMiddleware.cs index 513a6737d..cd9413e6d 100644 --- a/src/modules/Elsa.Workflows.Core/Pipelines/ActivityExecution/Components/ActivityInvokerMiddleware.cs +++ b/src/modules/Elsa.Workflows.Core/Pipelines/ActivityExecution/Components/ActivityInvokerMiddleware.cs @@ -36,13 +36,12 @@ public class ActivityInvokerMiddleware : IActivityExecutionMiddleware var executeDelegate = workflowExecution.ExecuteDelegate ?? (ExecuteActivityDelegate)Delegate.CreateDelegate(typeof(ExecuteActivityDelegate), activity, methodInfo); // Record executing event. - LogExecutionRecord(context, WorkflowExecutionLogEventNames.Executing); + _logger.LogTrace("Executing activity {ActivityId}", activity.Id); await executeDelegate(context); // Record executed event. - var payload = context.JournalData.Any() ? context.JournalData : default; - LogExecutionRecord(context, WorkflowExecutionLogEventNames.Executed, payload: payload); + _logger.LogTrace("Executed activity {ActivityId}", activity.Id); // Reset execute delegate. workflowExecution.ExecuteDelegate = null; @@ -57,6 +56,4 @@ public class ActivityInvokerMiddleware : IActivityExecutionMiddleware workflowExecution.RegisterBookmarks(context.Bookmarks); } } - - private void LogExecutionRecord(ActivityExecutionContext context, string eventName, string? message = default, string? source = default, object? payload = default) => context.AddExecutionLogEntry(eventName, message, source, payload); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Services/ISignalHandler.cs b/src/modules/Elsa.Workflows.Core/Services/ISignalHandler.cs index 70bed4ce0..3658b33d1 100644 --- a/src/modules/Elsa.Workflows.Core/Services/ISignalHandler.cs +++ b/src/modules/Elsa.Workflows.Core/Services/ISignalHandler.cs @@ -9,15 +9,28 @@ public interface ISignalHandler public class SignalContext { - public SignalContext(ActivityExecutionContext activityExecutionContext, ActivityExecutionContext sourceActivityExecutionContext, CancellationToken cancellationToken) + public SignalContext(ActivityExecutionContext receiverActivityExecutionContext, ActivityExecutionContext senderActivityExecutionContext, CancellationToken cancellationToken) { - ActivityExecutionContext = activityExecutionContext; - SourceActivityExecutionContext = sourceActivityExecutionContext; + ReceiverActivityExecutionContext = receiverActivityExecutionContext; + SenderActivityExecutionContext = senderActivityExecutionContext; CancellationToken = cancellationToken; } - public ActivityExecutionContext ActivityExecutionContext { get; init; } - public ActivityExecutionContext SourceActivityExecutionContext { get; init; } + /// + /// The receiving the signal. + /// + public ActivityExecutionContext ReceiverActivityExecutionContext { get; init; } + + /// + /// The sending the signal. + /// + public ActivityExecutionContext SenderActivityExecutionContext { get; init; } + + /// + /// Returns true if the receiver is the same as the sender. + /// + public bool IsSelf => SenderActivityExecutionContext.Activity == ReceiverActivityExecutionContext.Activity; + public CancellationToken CancellationToken { get; init; } internal bool StopPropagationRequested { get; private set; }