Refactor activity execution failure event publishing logic

This commit is contained in:
Sipke Schoorstra 2021-10-21 11:21:25 +02:00
parent bccb32b003
commit 7ee50cf6af
3 changed files with 34 additions and 31 deletions

View file

@ -4,9 +4,9 @@ using MediatR;
namespace Elsa.Events
{
public class ActivityExecutionResultFailed : INotification
public class ActivityExecutionFailed : INotification
{
public ActivityExecutionResultFailed(Exception exception, ActivityExecutionContext activityExecutionContext)
public ActivityExecutionFailed(Exception exception, ActivityExecutionContext activityExecutionContext)
{
Exception = exception;
ActivityExecutionContext = activityExecutionContext;

View file

@ -299,7 +299,7 @@ namespace Elsa.Services.Workflows
}
catch (Exception e)
{
await _mediator.Publish(new ActivityExecutionResultFailed(e, activityExecutionContext), cancellationToken);
await _mediator.Publish(new ActivityExecutionFailed(e, activityExecutionContext), cancellationToken);
throw;
}
}
@ -328,7 +328,6 @@ namespace Elsa.Services.Workflows
_logger.LogWarning(e, "Failed to run activity {ActivityId} of workflow {WorkflowInstanceId}", activity.Id, activityExecutionContext.WorkflowInstance.Id);
activityExecutionContext.Fault(e);
await _mediator.Publish(new ActivityFaulted(e, activityExecutionContext, activity), cancellationToken);
await _mediator.Publish(new ActivityExecutionResultFailed(e, activityExecutionContext), cancellationToken);
}
return null;

View file

@ -1,3 +1,4 @@
using System;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
@ -6,13 +7,14 @@ using Elsa.Events;
using Elsa.Models;
using Elsa.Server.Api.Models;
using Elsa.Server.Api.Services;
using Elsa.Services.Models;
using MediatR;
using Newtonsoft.Json;
using Newtonsoft.Json.Linq;
namespace Elsa.Server.Api.Handlers
{
public class ActivityExecutionResultExecutedHandler : INotificationHandler<ActivityExecutionResultExecuted>, INotificationHandler<ActivityExecutionResultFailed>
public class ActivityExecutionResultExecutedHandler : INotificationHandler<ActivityExecutionResultExecuted>, INotificationHandler<ActivityExecutionFailed>, INotificationHandler<ActivityFaulted>
{
private readonly IWorkflowTestService _workflowTestService;
@ -34,6 +36,7 @@ namespace Elsa.Server.Api.Handlers
};
var body = context.Input != null ? ((dynamic)context.Input).Body : null;
if (body != null)
data["Body"] = JToken.FromObject(body);
@ -55,6 +58,30 @@ namespace Elsa.Server.Api.Handlers
await _workflowTestService.DispatchMessage(signalRConnectionId, message);
}
public Task Handle(ActivityExecutionFailed notification, CancellationToken cancellationToken) => HandleFaultedExecutionAsync(notification.ActivityExecutionContext, notification.Exception);
public Task Handle(ActivityFaulted notification, CancellationToken cancellationToken) => HandleFaultedExecutionAsync(notification.ActivityExecutionContext, notification.Exception);
private async Task HandleFaultedExecutionAsync(ActivityExecutionContext context, Exception exception)
{
var signalRConnectionId = context.WorkflowExecutionContext.WorkflowInstance.GetMetadata("signalRConnectionId")?.ToString();
if (string.IsNullOrWhiteSpace(signalRConnectionId)) return;
var innerException = exception;
while (innerException?.InnerException != null) innerException = innerException.InnerException;
var message = new WorkflowTestMessage
{
WorkflowInstanceId = context.WorkflowInstance.Id,
CorrelationId = context.CorrelationId,
ActivityId = context.ActivityId,
Error = innerException?.ToString()
};
message.WorkflowStatus = message.Status = "Failed";
await _workflowTestService.DispatchMessage(signalRConnectionId, message);
}
private string GetExecutionResult(IActivityExecutionResult activityExecutionResult)
{
@ -68,40 +95,17 @@ namespace Elsa.Server.Api.Handlers
case DoneResult:
case OutcomeResult:
var outcomeResult = (OutcomeResult)activityExecutionResult;
status = outcomeResult.Outcomes.FirstOrDefault();
status = outcomeResult.Outcomes.First();
break;
case FaultResult:
status = "Failed";
break;
case CombinedResult:
var combinedResult = (CombinedResult)activityExecutionResult;
status = GetExecutionResult(combinedResult.Results.FirstOrDefault());
case CombinedResult combinedResult:
status = GetExecutionResult(combinedResult.Results.First());
break;
}
return status;
}
public async Task Handle(ActivityExecutionResultFailed notification, CancellationToken cancellationToken)
{
var context = notification.ActivityExecutionContext;
var signalRConnectionId = context.WorkflowExecutionContext.WorkflowInstance.GetMetadata("signalRConnectionId")?.ToString();
if (string.IsNullOrWhiteSpace(signalRConnectionId)) return;
var innerMostException = notification.Exception;
while (innerMostException.InnerException != null) innerMostException = innerMostException.InnerException;
var message = new WorkflowTestMessage
{
WorkflowInstanceId = context.WorkflowInstance.Id,
CorrelationId = context.CorrelationId,
ActivityId = context.ActivityId,
Error = innerMostException.ToString()
};
message.WorkflowStatus = message.Status = "Failed";
await _workflowTestService.DispatchMessage(signalRConnectionId, message);
}
}
}