From 7ee50cf6af824878941415c789e54bfcfec3fcc1 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 21 Oct 2021 11:21:25 +0200 Subject: [PATCH] Refactor activity execution failure event publishing logic --- ...ltFailed.cs => ActivityExecutionFailed.cs} | 4 +- .../Services/Workflows/WorkflowRunner.cs | 3 +- .../ActivityExecutionResultExecutedHandler.cs | 58 ++++++++++--------- 3 files changed, 34 insertions(+), 31 deletions(-) rename src/core/Elsa.Abstractions/Events/{ActivityExecutionResultFailed.cs => ActivityExecutionFailed.cs} (64%) diff --git a/src/core/Elsa.Abstractions/Events/ActivityExecutionResultFailed.cs b/src/core/Elsa.Abstractions/Events/ActivityExecutionFailed.cs similarity index 64% rename from src/core/Elsa.Abstractions/Events/ActivityExecutionResultFailed.cs rename to src/core/Elsa.Abstractions/Events/ActivityExecutionFailed.cs index 9d5d98dad..7b15bf4ed 100644 --- a/src/core/Elsa.Abstractions/Events/ActivityExecutionResultFailed.cs +++ b/src/core/Elsa.Abstractions/Events/ActivityExecutionFailed.cs @@ -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; diff --git a/src/core/Elsa.Core/Services/Workflows/WorkflowRunner.cs b/src/core/Elsa.Core/Services/Workflows/WorkflowRunner.cs index 0bf4fde6d..38049cf01 100644 --- a/src/core/Elsa.Core/Services/Workflows/WorkflowRunner.cs +++ b/src/core/Elsa.Core/Services/Workflows/WorkflowRunner.cs @@ -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; diff --git a/src/server/Elsa.Server.Api/Handlers/ActivityExecutionResultExecutedHandler.cs b/src/server/Elsa.Server.Api/Handlers/ActivityExecutionResultExecutedHandler.cs index 7d2498b0a..20b8aef55 100644 --- a/src/server/Elsa.Server.Api/Handlers/ActivityExecutionResultExecutedHandler.cs +++ b/src/server/Elsa.Server.Api/Handlers/ActivityExecutionResultExecutedHandler.cs @@ -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, INotificationHandler + public class ActivityExecutionResultExecutedHandler : INotificationHandler, INotificationHandler, INotificationHandler { 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); - } } } \ No newline at end of file