Remove signal capturing phase (#5501)

* Add logging to the WorkflowRunner service

The WorkflowRunner service now uses the Microsoft.Extensions.Logging namespace to log the workflow execution context. These changes include passing the ILogger<WorkflowRunner> logger dependency through the constructor and implementing the context logging functionality in the RunAsync method.

* Update error handling in FindActivityDescriptor method

FindActivityDescriptor method's error handling has been updated. Now, instead of throwing exception, it returns null when an activity descriptor can't be found. Also, a logger warning has been added to indicate when this situation occurs. This change helps avoiding unexpected disruptions and improving debugging experiences.

* Add ReSharper properties to .editorconfig

This commit introduces specific ReSharper properties to the .editorconfig file. This update will maintain a consistent configuration of ReSharper across different development environments, intending to improve coding standard consistency.

* Add thread and activity ID to Flowchart logging scope

The Flowchart activity in Elsa Workflows Core module has been modified to include thread and activity ID in its logging scope. This change will provide more granular information when debugging workflow execution. Additionally, the definition for outcomeNames has been streamlined.

* Add ActivityInstanceId to logger scope

Added "ActivityInstanceId" as part of the logging scope dictionary in the Flowchart module. The new key records the context's target context's id for improved debugging capabilities.

* Remove signal capturing functionality from workflow activities

Removed the functionality related to signal capturing from the Elsa workflow activities. This refactor involves changes in core classes such as Activity, Behavior, and Flowchart and removes associated methods and handlers. This simplifies the signal handling process by only allowing activities to receive signals, eliminating the previous two-step process of capturing and receiving.

* Update debugging messages in Flowchart.cs

Clarified the debugging message when there's an existing join context. Removed unnecessary logging for "No pending work found", "No faulted activities found", and "Completing flowchart". This will make the debugging log less cluttered and more focused on relevant information.

* Refactor logging messages in Flowchart activity

The commit removes the verbose logging message indicating the completion of a terminal activity in the flowchart Context. This logging message was unnecessary and was generating excessive log messages. The log message for new join activities was also updated to accurately reflect the creation of a new join context.

* Disable SonarCloud analysis from workflow

The SonarCloud analysis steps, including scan setup, run and end steps have been commented out in the GitHub workflow. This is a temporary change to speed up build times while troubleshooting an issue.

* Uncommented SonarCloud analysis related code in packages.yml

In this commit, the parts of the code related to the set up of JDK 17, SonarScanner for .NET, Coverlet for code coverage, and SonarCloud analysis were uncommented in the GitHub Actions workflow packages.yml file. This will enable those tools and services during the execution of the workflow, improving code quality and test coverage.

* Change position of root assignment comment in .editorconfig

* Format method signatures in WorkflowRunner

Changed the method signatures in the WorkflowRunner class to be in a single line for readability and to follow coding standards. The refactor involves three RunAsync method overloads, contributing to the overall cleanliness and readability of the source code.
This commit is contained in:
Sipke Schoorstra 2024-06-04 09:54:27 +02:00 committed by GitHub
parent 100ece8278
commit bd2e70cbe1
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
8 changed files with 49 additions and 201 deletions

View file

@ -228,3 +228,8 @@ dotnet_naming_style.begins_with_i.required_prefix = I
dotnet_naming_style.begins_with_i.required_suffix =
dotnet_naming_style.begins_with_i.word_separator =
dotnet_naming_style.begins_with_i.capitalization = pascal_case
# ReSharper properties
resharper_max_array_initializer_elements_on_line = 50
resharper_max_initializer_elements_on_line = 1
resharper_wrap_array_initializer_style = chop_if_long

View file

@ -17,7 +17,6 @@ namespace Elsa.Workflows;
public abstract class Activity : IActivity, ISignalHandler
{
private readonly ICollection<SignalHandlerRegistration> _signalReceivedHandlers = new List<SignalHandlerRegistration>();
private readonly ICollection<SignalHandlerRegistration> _signalCapturedHandlers = new List<SignalHandlerRegistration>();
/// <summary>
/// Constructor.
@ -161,44 +160,6 @@ public abstract class Activity : IActivity, ISignalHandler
return ValueTask.CompletedTask;
});
}
/// <summary>
/// Override this method to handle any signals sent from downstream activities.
/// </summary>
protected virtual ValueTask OnCaptureSignalAsync(object signal, SignalContext context)
{
OnSignalCaptured(signal, context);
return ValueTask.CompletedTask;
}
/// <summary>
/// Override this method to handle any signals sent from downstream activities.
/// </summary>
protected virtual void OnSignalCaptured(object signal, SignalContext context)
{
}
/// <summary>
/// Register a signal handler delegate.
/// </summary>
protected void OnSignalCaptured(Type signalType, Func<object, SignalContext, ValueTask> handler) => _signalCapturedHandlers.Add(new SignalHandlerRegistration(signalType, handler));
/// <summary>
/// Register a signal handler delegate.
/// </summary>
protected void OnSignalCaptured<T>(Func<T, SignalContext, ValueTask> handler) => OnSignalCaptured(typeof(T), (signal, context) => handler((T)signal, context));
/// <summary>
/// Register a signal handler delegate.
/// </summary>
protected void OnSignalCaptured<T>(Action<T, SignalContext> handler)
{
OnSignalCaptured<T>((signal, context) =>
{
handler(signal, context);
return ValueTask.CompletedTask;
});
}
/// <summary>
/// Notify the workflow that this activity completed.
@ -220,23 +181,7 @@ public abstract class Activity : IActivity, ISignalHandler
// Invoke behaviors.
foreach (var behavior in Behaviors) await behavior.ExecuteAsync(context);
}
async ValueTask ISignalHandler.CaptureSignalAsync(object signal, SignalContext context)
{
// Give derived activity a chance to do something with the signal.
await OnCaptureSignalAsync(signal, context);
// Invoke registered signal delegates for this particular type of signal.
var signalType = signal.GetType();
var handlers = _signalCapturedHandlers.Where(x => x.SignalType == signalType);
foreach (var registration in handlers)
await registration.Handler(signal, context);
// Invoke behaviors.
foreach (var behavior in Behaviors) await behavior.CaptureSignalAsync(signal, context);
}
async ValueTask ISignalHandler.ReceiveSignalAsync(object signal, SignalContext context)
{
// Give derived activity a chance to do something with the signal.

View file

@ -7,7 +7,6 @@ namespace Elsa.Workflows;
public abstract class Behavior : IBehavior
{
private readonly ICollection<SignalHandlerRegistration> _signalReceivedHandlers = new List<SignalHandlerRegistration>();
private readonly ICollection<SignalHandlerRegistration> _signalCapturedHandlers = new List<SignalHandlerRegistration>();
/// <summary>
/// Initializes a new instance of the <see cref="Behavior"/> class.
@ -71,54 +70,6 @@ public abstract class Behavior : IBehavior
{
}
/// <summary>
/// Registers a delegate to be invoked when a signal of the specified type is received.
/// </summary>
/// <param name="signalType">The type of signal to register a handler for.</param>
/// <param name="handler">The delegate to invoke when a signal of the specified type is received.</param>
protected void OnSignalCaptured(Type signalType, Func<object, SignalContext, ValueTask> handler) => _signalCapturedHandlers.Add(new SignalHandlerRegistration(signalType, handler));
/// <summary>
/// Registers a delegate to be invoked when a signal of the specified type is received.
/// </summary>
/// <param name="handler">The delegate to invoke when a signal of the specified type is received.</param>
/// <typeparam name="T">The type of signal to register a handler for.</typeparam>
protected void OnSignalCaptured<T>(Func<T, SignalContext, ValueTask> handler) => OnSignalCaptured(typeof(T), (signal, context) => handler((T)signal, context));
/// <summary>
/// Registers a delegate to be invoked when a signal of the specified type is received.
/// </summary>
/// <param name="handler">The delegate to invoke when a signal of the specified type is received.</param>
/// <typeparam name="T">The type of signal to register a handler for.</typeparam>
protected void OnSignalCaptured<T>(Action<T, SignalContext> handler)
{
OnSignalCaptured<T>((signal, context) =>
{
handler(signal, context);
return ValueTask.CompletedTask;
});
}
/// <summary>
/// Registers a delegate to be invoked when a signal of the specified type is received.
/// </summary>
/// <param name="signal">The type of signal to register a handler for.</param>
/// <param name="context">The signal context.</param>
protected virtual ValueTask OnSignalCapturedAsync(object signal, SignalContext context)
{
OnSignalCaptured(signal, context);
return ValueTask.CompletedTask;
}
/// <summary>
/// Registers a delegate to be invoked when a signal of the specified type is received.
/// </summary>
/// <param name="signal">The signal to register a handler for.</param>
/// <param name="context">The signal context.</param>
protected virtual void OnSignalCaptured(object signal, SignalContext context)
{
}
/// <summary>
///
/// </summary>
@ -137,20 +88,7 @@ public abstract class Behavior : IBehavior
protected virtual void Execute(ActivityExecutionContext context)
{
}
async ValueTask ISignalHandler.CaptureSignalAsync(object signal, SignalContext context)
{
// Give derived activity a chance to do something with the signal.
await OnSignalCapturedAsync(signal, context);
// Invoke registered signal delegates for this particular type of signal.
var signalType = signal.GetType();
var handlers = _signalCapturedHandlers.Where(x => x.SignalType == signalType);
foreach (var registration in handlers)
await registration.Handler(signal, context);
}
async ValueTask ISignalHandler.ReceiveSignalAsync(object signal, SignalContext context)
{
// Give derived activity a chance to do something with the signal.

View file

@ -20,7 +20,6 @@ namespace Elsa.Workflows.Activities.Flowchart.Activities;
public class Flowchart : Container
{
internal const string ScopeProperty = "Scope";
internal const string BranchMonitorsProperty = "BranchMonitoring";
/// <inheritdoc />
public Flowchart([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line)
@ -45,75 +44,51 @@ public class Flowchart : Container
/// <inheritdoc />
protected override async ValueTask ScheduleChildrenAsync(ActivityExecutionContext context)
{
var logger = context.GetRequiredService<ILogger<Flowchart>>();
var startActivity = GetStartActivity(context);
if (startActivity == null)
{
// Nothing else to execute.
logger.LogDebug("No start activity found. Completing flowchart");
await context.CompleteActivityAsync();
return;
}
// Schedule the start activity.
logger.LogDebug("Scheduling activity: {StartActivityId}", startActivity.Id);
await context.ScheduleActivityAsync(startActivity, OnChildCompletedAsync);
}
private IActivity? GetStartActivity(ActivityExecutionContext context)
{
var logger = context.GetRequiredService<ILogger<Flowchart>>();
logger.LogDebug("Looking for start activity...");
// If there's a trigger that triggered this workflow, use that.
var triggerActivityId = context.WorkflowExecutionContext.TriggerActivityId;
var triggerActivity = triggerActivityId != null ? Activities.FirstOrDefault(x => x.Id == triggerActivityId) : default;
if (triggerActivity != null)
{
logger.LogDebug("Found trigger activity: {TriggerActivityId}", triggerActivityId);
return triggerActivity;
}
// If an explicit Start activity was provided, use that.
if (Start != null)
{
logger.LogDebug("An explicit start activity was provided: {StartActivityId}", Start.Id);
return Start;
}
// If there is a Start activity on the flowchart, use that.
var startActivity = Activities.FirstOrDefault(x => x is Start);
if (startActivity != null)
{
logger.LogDebug("A Start activity was found: {StartActivityId}", startActivity.Id);
return startActivity;
}
// If there's an activity marked as "Can Start Workflow", use that.
var canStartWorkflowActivity = Activities.FirstOrDefault(x => x.GetCanStartWorkflow());
if (canStartWorkflowActivity != null)
{
logger.LogDebug("An activity marked as 'Can Start Workflow' was found: {CanStartWorkflowActivityId}", canStartWorkflowActivity.Id);
return canStartWorkflowActivity;
}
// If there is a single activity that has no inbound connections, use that.
var root = GetRootActivity();
if (root != null)
{
logger.LogDebug("Found a single activity with no inbound connections: {ActivityId}", root.Id);
return root;
}
// If no start activity found, return the first activity.
logger.LogDebug("No start activity found. Using the first activity");
return Activities.FirstOrDefault();
}
@ -164,21 +139,24 @@ public class Flowchart : Container
private async ValueTask OnChildCompletedAsync(ActivityCompletedContext context)
{
var logger = context.GetRequiredService<ILogger<Flowchart>>();
var loggerScopeState = new Dictionary<string, object>
{
["ThreadId"] = Thread.CurrentThread.ManagedThreadId,
["ActivityId"] = Id,
["ActivityInstanceId"] = context.TargetContext.Id
};
using var loggerScope = logger.BeginScope(loggerScopeState);
var flowchartContext = context.TargetContext;
var completedActivityContext = context.ChildContext;
var completedActivity = completedActivityContext.Activity;
var result = context.Result;
logger.LogDebug("Child activity {ActivityId} completed with status {ActivityStatus}", completedActivity.Id, completedActivityContext.Status);
// If the complete activity's status is anything but "Completed", do not schedule its outbound activities.
var scheduleChildren = completedActivityContext.Status == ActivityStatus.Completed;
var outcomeNames = result is Outcomes outcomes
? outcomes.Names
: new[]
{
default(string), "Done"
};
: [null!, "Done"];
// Only query the outbound connections if the completed activity wasn't already completed.
var outboundConnections = Connections.Where(connection => connection.Source.Activity == completedActivity && outcomeNames.Contains(connection.Source.Port)).ToList();
@ -190,18 +168,12 @@ public class Flowchart : Container
// If the complete activity is a terminal node, complete the flowchart immediately.
if (completedActivity is ITerminalNode)
{
logger.LogDebug("Completed activity {ActivityId} is a terminal activity. Completing flowchart", completedActivity.Id);
await flowchartContext.CompleteActivityAsync();
}
else if (scheduleChildren)
{
if (children.Any())
{
if (children.Count == 1)
logger.LogDebug("Found 1 child for activity {ActivityId}: {ChildActivityId}", completedActivity.Id, children.First().Id);
else
logger.LogDebug("Found {Count} children for activity {ActivityId}: {ChildActivityIds}", children.Count, completedActivity.Id, children.Select(x => x.Id).ToList());
scope.AddActivities(children);
// Schedule each child, but only if all of its left inbound activities have already executed.
@ -212,7 +184,6 @@ public class Flowchart : Container
// If the completed activity is not part of the left inbound path, always allow its children to be scheduled.
if (!inboundActivities.Contains(completedActivity))
{
logger.LogDebug("Scheduling child activity {ChildActivityId}", activity.Id);
await flowchartContext.ScheduleActivityAsync(activity, OnChildCompletedAsync);
continue;
}
@ -225,7 +196,6 @@ public class Flowchart : Container
if (haveInboundActivitiesExecuted)
{
logger.LogDebug("Scheduling child activity {ChildActivityId}", activity.Id);
await flowchartContext.ScheduleActivityAsync(activity, OnChildCompletedAsync);
}
}
@ -241,11 +211,10 @@ public class Flowchart : Container
};
if (joinContext != null)
logger.LogDebug("Next activity {ChildActivityId} is a join activity. Attaching to existing context {JoinContext}", activity.Id, joinContext.Id);
logger.LogDebug("Next activity {ChildActivityId} is a join activity. Attaching to existing join context {JoinContext}", activity.Id, joinContext.Id);
else
logger.LogDebug("Next activity {ChildActivityId} is a join activity", activity.Id);
logger.LogDebug("Scheduling child activity {ChildActivityId}", activity.Id);
logger.LogDebug("Next activity {ChildActivityId} is a join activity. Creating new join context", activity.Id);
await flowchartContext.ScheduleActivityAsync(activity, scheduleWorkOptions);
}
}
@ -253,7 +222,6 @@ public class Flowchart : Container
if (!children.Any())
{
logger.LogDebug("No children found for activity {ActivityId}", completedActivity.Id);
await CompleteIfNoPendingWorkAsync(flowchartContext);
}
}
@ -268,13 +236,10 @@ public class Flowchart : Container
if (!hasPendingWork)
{
logger.LogDebug("No pending work found");
var hasFaultedActivities = context.GetActiveChildren().Any(x => x.Status == ActivityStatus.Faulted);
if (!hasFaultedActivities)
{
logger.LogDebug("No faulted activities found");
logger.LogDebug("Completing flowchart");
await context.CompleteActivityAsync();
}
}

View file

@ -5,11 +5,6 @@ namespace Elsa.Workflows.Contracts;
/// </summary>
public interface ISignalHandler
{
/// <summary>
/// Captures a signal.
/// </summary>
ValueTask CaptureSignalAsync(object signal, SignalContext context);
/// <summary>
/// Receives a signal.
/// </summary>

View file

@ -363,26 +363,9 @@ public static class ActivityExecutionContextExtensions
public static async ValueTask SendSignalAsync(this ActivityExecutionContext context, object signal)
{
var receivingContexts = new[] { context }.Concat(context.GetAncestors()).ToList();
var capturingContexts = receivingContexts.AsEnumerable().Reverse().ToList();
var logger = context.GetRequiredService<ILogger<ActivityExecutionContext>>();
// Let all ancestors capture the signal.
foreach (var ancestorContext in capturingContexts)
{
var signalContext = new SignalContext(ancestorContext, context, context.CancellationToken);
if (ancestorContext.Activity is not ISignalHandler handler)
continue;
logger.LogDebug("Capturing signal {SignalType} on activity {ActivityId} of type {ActivityType}", signal.GetType().Name, ancestorContext.Activity.Id, ancestorContext.Activity.Type);
await handler.CaptureSignalAsync(signal, signalContext);
if (signalContext.StopPropagationRequested)
{
logger.LogDebug("Propagation of signal {SignalType} on activity {ActivityId} of type {ActivityType} was stopped", signal.GetType().Name, ancestorContext.Activity.Id, ancestorContext.Activity.Type);
return;
}
}
var signalType = signal.GetType();
var signalTypeName = signalType.Name;
// Let all ancestors receive the signal.
foreach (var ancestorContext in receivingContexts)
@ -392,12 +375,12 @@ public static class ActivityExecutionContextExtensions
if (ancestorContext.Activity is not ISignalHandler handler)
continue;
logger.LogDebug("Receiving signal {SignalType} on activity {ActivityId} of type {ActivityType}", signal.GetType().Name, ancestorContext.Activity.Id, ancestorContext.Activity.Type);
logger.LogDebug("Receiving signal {SignalType} on activity {ActivityId} of type {ActivityType}", signalTypeName, ancestorContext.Activity.Id, ancestorContext.Activity.Type);
await handler.ReceiveSignalAsync(signal, signalContext);
if (signalContext.StopPropagationRequested)
{
logger.LogDebug("Propagation of signal {SignalType} on activity {ActivityId} of type {ActivityType} was stopped", signal.GetType().Name, ancestorContext.Activity.Id, ancestorContext.Activity.Type);
logger.LogDebug("Propagation of signal {SignalType} on activity {ActivityId} of type {ActivityType} was stopped", signalTypeName, ancestorContext.Activity.Id, ancestorContext.Activity.Type);
return;
}
}

View file

@ -6,6 +6,7 @@ using Elsa.Workflows.Models;
using Elsa.Workflows.Notifications;
using Elsa.Workflows.Options;
using Elsa.Workflows.State;
using Microsoft.Extensions.Logging;
namespace Elsa.Workflows.Services;
@ -17,7 +18,8 @@ public class WorkflowRunner(
IWorkflowBuilderFactory workflowBuilderFactory,
IWorkflowGraphBuilder workflowGraphBuilder,
IIdentityGenerator identityGenerator,
INotificationSender notificationSender)
INotificationSender notificationSender,
ILogger<WorkflowRunner> logger)
: IWorkflowRunner
{
/// <inheritdoc />
@ -107,7 +109,7 @@ public class WorkflowRunner(
var workflowGraph = await workflowGraphBuilder.BuildAsync(workflow, cancellationToken);
return await RunAsync(workflowGraph, workflowState, options, cancellationToken);
}
/// <inheritdoc />
public async Task<RunWorkflowResult> RunAsync(WorkflowGraph workflowGraph, WorkflowState workflowState, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default)
{
@ -161,7 +163,8 @@ public class WorkflowRunner(
}
else if (activityInstanceId != null)
{
var activityExecutionContext = workflowExecutionContext.ActivityExecutionContexts.FirstOrDefault(x => x.Id == activityInstanceId) ?? throw new Exception("No activity execution context found with the specified ID.");
var activityExecutionContext = workflowExecutionContext.ActivityExecutionContexts.FirstOrDefault(x => x.Id == activityInstanceId) ??
throw new Exception("No activity execution context found with the specified ID.");
workflowExecutionContext.ScheduleActivityExecutionContext(activityExecutionContext);
}
else if (workflowExecutionContext.Scheduler.HasAny)
@ -180,6 +183,12 @@ public class WorkflowRunner(
/// <inheritdoc />
public async Task<RunWorkflowResult> RunAsync(WorkflowExecutionContext workflowExecutionContext)
{
var workflowInstanceId = workflowExecutionContext.Id;
var logContext = new Dictionary<string, object>
{
["WorkflowInstanceId"] = workflowInstanceId
};
using var loggingScope = logger.BeginScope(logContext);
var workflow = workflowExecutionContext.Workflow;
var applicationCancellationToken = workflowExecutionContext.CancellationTokens.ApplicationCancellationToken;
var systemCancellationToken = workflowExecutionContext.CancellationTokens.SystemCancellationToken;

View file

@ -165,10 +165,10 @@ public class WorkflowDefinitionActivity : Composite, IInitializable
return workflowGraph;
}
private ActivityDescriptor FindActivityDescriptor(IServiceProvider serviceProvider)
private ActivityDescriptor? FindActivityDescriptor(IServiceProvider serviceProvider)
{
var activityRegistry = serviceProvider.GetRequiredService<IActivityRegistry>();
return activityRegistry.Find(Type, Version) ?? activityRegistry.Find(Type) ?? throw new Exception($"Could not find activity descriptor for {Type}.");
return activityRegistry.Find(Type, Version) ?? activityRegistry.Find(Type);
}
async ValueTask IInitializable.InitializeAsync(InitializationContext context)
@ -188,10 +188,18 @@ public class WorkflowDefinitionActivity : Composite, IInitializable
var activityDescriptor = FindActivityDescriptor(serviceProvider);
// Declare input and output variables.
DeclareInputAsVariables(activityDescriptor, (_, variable) => Variables.Declare(variable));
DeclareOutputAsVariables(activityDescriptor, (_, variable) => Variables.Declare(variable));
if (activityDescriptor == null)
{
var logger = serviceProvider.GetRequiredService<ILogger<WorkflowDefinitionActivity>>();
logger.LogWarning("Could not find activity descriptor for activity type {ActivityType}", Type);
}
else
{
// Declare input and output variables.
DeclareInputAsVariables(activityDescriptor, (_, variable) => Variables.Declare(variable));
DeclareOutputAsVariables(activityDescriptor, (_, variable) => Variables.Declare(variable));
}
// Set the root activity.
Root = workflowGraph.Workflow;
}