Implement Started and Completed execution events

This commit is contained in:
Sipke Schoorstra 2022-06-03 14:22:55 +02:00
parent 4b714a5021
commit 5a2b539637
9 changed files with 62 additions and 27 deletions

View file

@ -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)

View file

@ -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();
}
}

View file

@ -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();
}
}

View file

@ -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;
/// <summary>
/// Records Start and Completed entries of the owning activity.
/// </summary>
public class ExecutionLoggingBehavior : Behavior
{
public ExecutionLoggingBehavior(IActivity owner) : base(owner)
{
OnSignalReceived<ActivityCompleted>(OnActivityCompleted);
}
protected override void Execute(ActivityExecutionContext context)
{
context.AddExecutionLogEntry("Started");
}
private void OnActivityCompleted(ActivityCompleted signal, SignalContext context)
{
if (context.IsSelf)
context.SenderActivityExecutionContext.AddExecutionLogEntry("Completed");
}
}

View file

@ -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);

View file

@ -13,9 +13,6 @@ public static class ActivityExecutionContextExtensions
public static T GetInput<T>(this ActivityExecutionContext context) => context.GetInput<T>(typeof(T).Name);
public static T GetInput<T>(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<T?> EvaluateAsync<T>(this ActivityExecutionContext context, Input<T> input)
{
var evaluator = context.GetRequiredService<IExpressionEvaluator>();
@ -87,7 +84,7 @@ public static class ActivityExecutionContextExtensions
locationReference.Set(context, value);
return value;
}
/// <summary>
/// Returns a flattened list of the current context's ancestors.
/// </summary>
@ -107,7 +104,7 @@ public static class ActivityExecutionContextExtensions
/// Returns a flattened list of the current context's immediate children.
/// </summary>
/// <returns></returns>
public static IEnumerable<ActivityExecutionContext> GetChildren(this ActivityExecutionContext context) =>
public static IEnumerable<ActivityExecutionContext> GetChildren(this ActivityExecutionContext context) =>
context.WorkflowExecutionContext.ActivityExecutionContexts.Where(x => x.ParentActivityExecutionContext == context);
/// <summary>
@ -118,13 +115,13 @@ public static class ActivityExecutionContextExtensions
// Detach child activity execution contexts.
context.WorkflowExecutionContext.RemoveActivityExecutionContexts(context.GetChildren());
}
/// <summary>
/// Send a signal up the current branch.
/// </summary>
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)

View file

@ -14,9 +14,9 @@ public abstract class Activity : IActivity, ISignalHandler
protected Activity()
{
TypeName = ActivityTypeNameHelper.GenerateTypeName(GetType());
Behaviors.Add<ExecutionLoggingBehavior>(this);
Behaviors.Add<ScheduledChildCallbackBehavior>(this);
Behaviors.Add<AutoCompleteBehavior>(this);
Behaviors.Add<CreateCompletedLogRecord>(this);
}
protected Activity(string activityType) : this()

View file

@ -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);
}

View file

@ -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; }
/// <summary>
/// The <see cref="ActivityExecutionContext"/> receiving the signal.
/// </summary>
public ActivityExecutionContext ReceiverActivityExecutionContext { get; init; }
/// <summary>
/// The <see cref="ActivityExecutionContext"/> sending the signal.
/// </summary>
public ActivityExecutionContext SenderActivityExecutionContext { get; init; }
/// <summary>
/// Returns true if the receiver is the same as the sender.
/// </summary>
public bool IsSelf => SenderActivityExecutionContext.Activity == ReceiverActivityExecutionContext.Activity;
public CancellationToken CancellationToken { get; init; }
internal bool StopPropagationRequested { get; private set; }