Implement workflow-level exception handling

Fixes #4451
This commit is contained in:
Sipke Schoorstra 2023-09-18 20:52:53 +02:00
parent e52c4eb91f
commit fbdd555dd9
6 changed files with 100 additions and 11 deletions

View file

@ -1,3 +1,4 @@
using Elsa.Common.Contracts;
using Elsa.Workflows.Core;
using Elsa.Workflows.Core.Contracts;
using Elsa.Workflows.Core.Models;
@ -109,4 +110,36 @@ public static class WorkflowExecutionContextExtensions
/// Returns true if all activities have completed or canceled, false otherwise.
/// </summary>
public static bool AllActivitiesCompleted(this WorkflowExecutionContext workflowExecutionContext) => workflowExecutionContext.ActiveActivityExecutionContexts.All(x => x.Status != ActivityStatus.Running);
/// <summary>
/// Adds a new <see cref="WorkflowExecutionLogEntry"/> to the execution log of the current <see cref="WorkflowExecutionContext"/>.
/// </summary>
/// <param name="context">The <see cref="WorkflowExecutionContext"/></param> being extended.
/// <param name="eventName">The name of the event.</param>
/// <param name="message">The message of the event.</param>
/// <param name="payload">Any contextual data related to this event.</param>
/// <returns>Returns the created <see cref="WorkflowExecutionLogEntry"/>.</returns>
public static WorkflowExecutionLogEntry AddExecutionLogEntry(this WorkflowExecutionContext context, string eventName, string? message = default, object? payload = default)
{
var now = context.GetRequiredService<ISystemClock>().UtcNow;
var logEntry = new WorkflowExecutionLogEntry(
context.Id,
default,
context.Workflow.Id,
context.Workflow.Type,
context.Workflow.Identity.Version,
context.Workflow.Name,
context.Workflow.Identity.Id,
default,
now,
context.ExecutionLogSequence++,
eventName,
message,
context.Workflow.GetSource(),
payload);
context.ExecutionLog.Add(logEntry);
return logEntry;
}
}

View file

@ -47,7 +47,9 @@ public class WorkflowsFeature : FeatureBase
/// <summary>
/// A delegate to configure the <see cref="IWorkflowExecutionPipeline"/>.
/// </summary>
public Action<IWorkflowExecutionPipelineBuilder> WorkflowExecutionPipeline { get; set; } = builder => builder.UseDefaultActivityScheduler();
public Action<IWorkflowExecutionPipelineBuilder> WorkflowExecutionPipeline { get; set; } = builder => builder
.UseExceptionHandling()
.UseDefaultActivityScheduler();
/// <summary>
/// A delegate to configure the <see cref="IActivityExecutionPipeline"/>.

View file

@ -0,0 +1,60 @@
using Elsa.Common.Contracts;
using Elsa.Extensions;
using Elsa.Workflows.Core.Contracts;
using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.Pipelines.WorkflowExecution;
using Elsa.Workflows.Core.State;
using Microsoft.Extensions.Logging;
namespace Elsa.Workflows.Core.Middleware.Workflows;
/// <summary>
/// Adds extension methods to <see cref="ExceptionHandlingMiddleware"/>.
/// </summary>
public static class ExceptionHandlingMiddlewareExtensions
{
/// <summary>
/// Installs the <see cref="ExceptionHandlingMiddleware"/> component in the activity execution pipeline.
/// </summary>
public static IWorkflowExecutionPipelineBuilder UseExceptionHandling(this IWorkflowExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder.UseMiddleware<ExceptionHandlingMiddleware>();
}
/// <summary>
/// Catches any exceptions thrown by downstream components and transitions the workflow into the faulted state.
/// </summary>
public class ExceptionHandlingMiddleware : IWorkflowExecutionMiddleware
{
private readonly WorkflowMiddlewareDelegate _next;
private readonly ISystemClock _systemClock;
private readonly ILogger<ExceptionHandlingMiddleware> _logger;
/// <summary>
/// Constructor.
/// </summary>
public ExceptionHandlingMiddleware(WorkflowMiddlewareDelegate next, ISystemClock systemClock, ILogger<ExceptionHandlingMiddleware> logger)
{
_next = next;
_systemClock = systemClock;
_logger = logger;
}
/// <inheritdoc />
public async ValueTask InvokeAsync(WorkflowExecutionContext context)
{
try
{
await _next(context);
}
catch (Exception e)
{
_logger.LogWarning(e, "An exception was caught from a downstream middleware component");
var exceptionState = ExceptionState.FromException(e);
var now = _systemClock.UtcNow;
var activity = context.Workflow;
var incident = new ActivityIncident(activity.Id, activity.Type, e.Message, exceptionState, now);
context.Incidents.Add(incident);
context.TransitionTo(WorkflowSubStatus.Faulted);
context.AddExecutionLogEntry("Faulted", e.Message, exceptionState);
}
}
}

View file

@ -1,9 +0,0 @@
namespace Elsa.Workflows.Core.Models;
/// <summary>
/// Holds information about a workflow fault.
/// </summary>
/// <param name="Exception">The exception that occurred</param>
/// <param name="Message">A description about the fault. Usually the exception message, if there waa an exception.</param>
/// <param name="FaultedActivityId">The ID of the activity that caused the workflow to fault.</param>
public record WorkflowFault(Exception? Exception, string Message, string? FaultedActivityId);

View file

@ -33,5 +33,7 @@ public class WorkflowExecutionPipeline : IWorkflowExecutionPipeline
/// <inheritdoc />
public async Task ExecuteAsync(WorkflowExecutionContext context) => await Pipeline(context);
private WorkflowMiddlewareDelegate CreateDefaultPipeline() => Setup(x => x.UseDefaultActivityScheduler());
private WorkflowMiddlewareDelegate CreateDefaultPipeline() => Setup(x => x
.UseExceptionHandling()
.UseDefaultActivityScheduler());
}

View file

@ -21,6 +21,7 @@ public static class WorkflowExecutionPipelineBuilderExtensions
.UseBookmarkPersistence()
.UseActivityExecutionLogPersistence()
.UseWorkflowExecutionLogPersistence()
.UseExceptionHandling()
.UseDefaultActivityScheduler();
/// <summary>