From 3c8d7c3c0da80ea8d4d9adfeacc80e7dd521e3b3 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sat, 12 Dec 2020 19:41:45 +0100 Subject: [PATCH] Extract workflow context handling to event handlers --- .../Events/WorkflowExecuted.cs | 2 +- .../Events/WorkflowExecuting.cs | 14 +++++ .../Handlers/PersistWorkflowContext.cs | 46 +++++++++++++++++ .../Handlers/RefreshWorkflowContext.cs | 51 +++++++++++++++++++ src/core/Elsa.Core/Services/WorkflowRunner.cs | 39 +------------- 5 files changed, 114 insertions(+), 38 deletions(-) create mode 100644 src/core/Elsa.Abstractions/Events/WorkflowExecuting.cs create mode 100644 src/core/Elsa.Core/Handlers/PersistWorkflowContext.cs create mode 100644 src/core/Elsa.Core/Handlers/RefreshWorkflowContext.cs diff --git a/src/core/Elsa.Abstractions/Events/WorkflowExecuted.cs b/src/core/Elsa.Abstractions/Events/WorkflowExecuted.cs index 15aa3cc10..d1a9e60b4 100644 --- a/src/core/Elsa.Abstractions/Events/WorkflowExecuted.cs +++ b/src/core/Elsa.Abstractions/Events/WorkflowExecuted.cs @@ -5,7 +5,7 @@ namespace Elsa.Events /// /// Published when a burst of execution finished. /// - public class WorkflowExecuted : WorkflowNotification + public class WorkflowExecuted : WorkflowNotification { public WorkflowExecuted(WorkflowExecutionContext workflowExecutionContext) : base(workflowExecutionContext) { diff --git a/src/core/Elsa.Abstractions/Events/WorkflowExecuting.cs b/src/core/Elsa.Abstractions/Events/WorkflowExecuting.cs new file mode 100644 index 000000000..dc81ca672 --- /dev/null +++ b/src/core/Elsa.Abstractions/Events/WorkflowExecuting.cs @@ -0,0 +1,14 @@ +using Elsa.Services.Models; + +namespace Elsa.Events +{ + /// + /// Published when a workflow is about to be executed. + /// + public class WorkflowExecuting : WorkflowNotification + { + public WorkflowExecuting(WorkflowExecutionContext workflowExecutionContext) : base(workflowExecutionContext) + { + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Handlers/PersistWorkflowContext.cs b/src/core/Elsa.Core/Handlers/PersistWorkflowContext.cs new file mode 100644 index 000000000..92210fb82 --- /dev/null +++ b/src/core/Elsa.Core/Handlers/PersistWorkflowContext.cs @@ -0,0 +1,46 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Events; +using Elsa.Models; +using Elsa.Services; +using Elsa.Services.Models; +using MediatR; + +namespace Elsa.Handlers +{ + public class PersistWorkflowContext : INotificationHandler, INotificationHandler + { + private readonly IWorkflowContextManager _workflowContextManager; + + public PersistWorkflowContext(IWorkflowContextManager workflowContextManager) + { + _workflowContextManager = workflowContextManager; + } + + + public async Task Handle(WorkflowExecuted notification, CancellationToken cancellationToken) + { + var workflowExecutionContext = notification.WorkflowExecutionContext; + var workflowInstance = workflowExecutionContext.WorkflowInstance; + workflowInstance.ContextId = await SaveWorkflowContextAsync(workflowExecutionContext, WorkflowContextFidelity.Burst, false, cancellationToken); + } + + public async Task Handle(ActivityExecuted notification, CancellationToken cancellationToken) + { + var workflowExecutionContext = notification.WorkflowExecutionContext; + var activityBlueprint = notification.Activity; + workflowExecutionContext.WorkflowInstance.ContextId = await SaveWorkflowContextAsync(workflowExecutionContext, WorkflowContextFidelity.Activity, activityBlueprint.SaveWorkflowContext, cancellationToken); + } + + private async ValueTask SaveWorkflowContextAsync(WorkflowExecutionContext workflowExecutionContext, WorkflowContextFidelity fidelity, bool always, CancellationToken cancellationToken) + { + var workflowContext = workflowExecutionContext.WorkflowContext; + + if (!always && (workflowContext == null || workflowExecutionContext.WorkflowBlueprint.ContextOptions?.ContextFidelity != fidelity)) + return workflowExecutionContext.WorkflowInstance.ContextId; + + var context = new SaveWorkflowContext(workflowExecutionContext); + return await _workflowContextManager.SaveContextAsync(context, cancellationToken); + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Handlers/RefreshWorkflowContext.cs b/src/core/Elsa.Core/Handlers/RefreshWorkflowContext.cs new file mode 100644 index 000000000..ea6b67edb --- /dev/null +++ b/src/core/Elsa.Core/Handlers/RefreshWorkflowContext.cs @@ -0,0 +1,51 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Events; +using Elsa.Models; +using Elsa.Services; +using Elsa.Services.Models; +using MediatR; + +namespace Elsa.Handlers +{ + public class RefreshWorkflowContext : INotificationHandler, INotificationHandler + { + private readonly IWorkflowContextManager _workflowContextManager; + + public RefreshWorkflowContext(IWorkflowContextManager workflowContextManager) + { + _workflowContextManager = workflowContextManager; + } + + public async Task Handle(WorkflowExecuting notification, CancellationToken cancellationToken) + { + var workflowExecutionContext = notification.WorkflowExecutionContext; + workflowExecutionContext.WorkflowContext = await LoadWorkflowContextAsync(workflowExecutionContext, WorkflowContextFidelity.Burst, false, cancellationToken); + } + + public async Task Handle(ActivityExecuting notification, CancellationToken cancellationToken) + { + var workflowExecutionContext = notification.WorkflowExecutionContext; + var workflowBlueprint = workflowExecutionContext.WorkflowBlueprint; + var activityBlueprint = notification.Activity; + + if (workflowBlueprint.ContextOptions?.ContextFidelity == WorkflowContextFidelity.Activity || activityBlueprint.LoadWorkflowContext || workflowExecutionContext.ContextHasChanged) + { + workflowExecutionContext.WorkflowContext = await LoadWorkflowContextAsync(workflowExecutionContext, WorkflowContextFidelity.Activity, activityBlueprint.LoadWorkflowContext || workflowExecutionContext.ContextHasChanged, cancellationToken); + workflowExecutionContext.ContextHasChanged = false; + } + } + + private async ValueTask LoadWorkflowContextAsync(WorkflowExecutionContext workflowExecutionContext, WorkflowContextFidelity fidelity, bool always, CancellationToken cancellationToken) + { + var workflowInstance = workflowExecutionContext.WorkflowInstance; + var workflowBlueprint = workflowExecutionContext.WorkflowBlueprint; + + if (!always && (workflowInstance.ContextId == null || workflowBlueprint.ContextOptions == null || workflowBlueprint.ContextOptions.ContextFidelity != fidelity)) + return null; + + var context = new LoadWorkflowContext(workflowExecutionContext); + return await _workflowContextManager.LoadContext(context, cancellationToken); + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/WorkflowRunner.cs b/src/core/Elsa.Core/Services/WorkflowRunner.cs index 2ec2a121b..52f8416ba 100644 --- a/src/core/Elsa.Core/Services/WorkflowRunner.cs +++ b/src/core/Elsa.Core/Services/WorkflowRunner.cs @@ -30,7 +30,6 @@ namespace Elsa.Services private readonly IWorkflowSelector _workflowSelector; private readonly IWorkflowInstanceStore _workflowInstanceManager; private readonly Func _workflowBuilderFactory; - private readonly IWorkflowContextManager _workflowContextManager; private readonly IMediator _mediator; private readonly IServiceProvider _serviceProvider; private readonly ILogger _logger; @@ -40,7 +39,6 @@ namespace Elsa.Services IWorkflowFactory workflowFactory, IWorkflowSelector workflowSelector, Func workflowBuilderFactory, - IWorkflowContextManager workflowContextManager, IMediator mediator, IServiceProvider serviceProvider, ILogger logger, IWorkflowInstanceStore workflowInstanceStore) @@ -48,7 +46,6 @@ namespace Elsa.Services _workflowRegistry = workflowRegistry; _workflowFactory = workflowFactory; _workflowBuilderFactory = workflowBuilderFactory; - _workflowContextManager = workflowContextManager; _mediator = mediator; _serviceProvider = serviceProvider; _logger = logger; @@ -160,7 +157,7 @@ namespace Elsa.Services { var workflowExecutionScope = _serviceProvider.CreateScope(); var workflowExecutionContext = new WorkflowExecutionContext(workflowExecutionScope, workflowBlueprint, workflowInstance, input); - workflowExecutionContext.WorkflowContext = await LoadWorkflowContextAsync(workflowExecutionContext, WorkflowContextFidelity.Burst, false, cancellationToken); + await _mediator.Publish(new WorkflowExecuting(workflowExecutionContext), cancellationToken); var activity = activityId != null ? workflowBlueprint.GetActivity(activityId) : default; @@ -179,7 +176,6 @@ namespace Elsa.Services break; } - workflowInstance.ContextId = await SaveWorkflowContextAsync(workflowExecutionContext, WorkflowContextFidelity.Burst, false, cancellationToken); await _mediator.Publish(new WorkflowExecuted(workflowExecutionContext), cancellationToken); var statusEvent = workflowExecutionContext.Status switch @@ -197,29 +193,6 @@ namespace Elsa.Services return workflowExecutionContext.WorkflowInstance; } - private async ValueTask LoadWorkflowContextAsync(WorkflowExecutionContext workflowExecutionContext, WorkflowContextFidelity fidelity, bool always, CancellationToken cancellationToken) - { - var workflowInstance = workflowExecutionContext.WorkflowInstance; - var workflowBlueprint = workflowExecutionContext.WorkflowBlueprint; - - if (!always && (workflowInstance.ContextId == null || workflowBlueprint.ContextOptions == null || workflowBlueprint.ContextOptions.ContextFidelity != fidelity)) - return null; - - var context = new LoadWorkflowContext(workflowExecutionContext); - return await _workflowContextManager.LoadContext(context, cancellationToken); - } - - private async ValueTask SaveWorkflowContextAsync(WorkflowExecutionContext workflowExecutionContext, WorkflowContextFidelity fidelity, bool always, CancellationToken cancellationToken) - { - var workflowContext = workflowExecutionContext.WorkflowContext; - - if (!always && (workflowContext == null || workflowExecutionContext.WorkflowBlueprint.ContextOptions?.ContextFidelity != fidelity)) - return workflowExecutionContext.WorkflowInstance.ContextId; - - var context = new SaveWorkflowContext(workflowExecutionContext); - return await _workflowContextManager.SaveContextAsync(context, cancellationToken); - } - private async Task BeginWorkflow(WorkflowExecutionContext workflowExecutionContext, IActivityBlueprint? activity, object? input, CancellationToken cancellationToken) { if (activity == null) @@ -270,21 +243,13 @@ namespace Elsa.Services var scheduledActivity = workflowExecutionContext.PopScheduledActivity(); var currentActivityId = scheduledActivity.ActivityId; var activityBlueprint = workflowBlueprint.GetActivity(currentActivityId)!; - - if (workflowBlueprint.ContextOptions?.ContextFidelity == WorkflowContextFidelity.Activity || activityBlueprint.LoadWorkflowContext || workflowExecutionContext.ContextHasChanged) - { - workflowExecutionContext.WorkflowContext = await LoadWorkflowContextAsync(workflowExecutionContext, WorkflowContextFidelity.Activity, activityBlueprint.LoadWorkflowContext || workflowExecutionContext.ContextHasChanged, cancellationToken); - workflowExecutionContext.ContextHasChanged = false; - } - var activityExecutionContext = new ActivityExecutionContext(scope, workflowExecutionContext, activityBlueprint, scheduledActivity.Input, cancellationToken); var activity = await activityExecutionContext.ActivateActivityAsync(cancellationToken); var result = await activityOperation(activityExecutionContext, activity); await _mediator.Publish(new ActivityExecuting(activityExecutionContext), cancellationToken); await result.ExecuteAsync(activityExecutionContext, cancellationToken); - await _mediator.Publish(new ActivityExecuted(activityExecutionContext), cancellationToken); workflowExecutionContext.WorkflowInstance.Output = activityExecutionContext.Output; - workflowExecutionContext.WorkflowInstance.ContextId = await SaveWorkflowContextAsync(workflowExecutionContext, WorkflowContextFidelity.Activity, activityBlueprint.SaveWorkflowContext, cancellationToken); + await _mediator.Publish(new ActivityExecuted(activityExecutionContext), cancellationToken); activityOperation = Execute; workflowExecutionContext.CompletePass();