Extract workflow context handling to event handlers
This commit is contained in:
parent
8dab2adca3
commit
3c8d7c3c0d
|
|
@ -5,7 +5,7 @@ namespace Elsa.Events
|
|||
/// <summary>
|
||||
/// Published when a burst of execution finished.
|
||||
/// </summary>
|
||||
public class WorkflowExecuted : WorkflowNotification
|
||||
public class WorkflowExecuted : WorkflowNotification
|
||||
{
|
||||
public WorkflowExecuted(WorkflowExecutionContext workflowExecutionContext) : base(workflowExecutionContext)
|
||||
{
|
||||
|
|
|
|||
14
src/core/Elsa.Abstractions/Events/WorkflowExecuting.cs
Normal file
14
src/core/Elsa.Abstractions/Events/WorkflowExecuting.cs
Normal file
|
|
@ -0,0 +1,14 @@
|
|||
using Elsa.Services.Models;
|
||||
|
||||
namespace Elsa.Events
|
||||
{
|
||||
/// <summary>
|
||||
/// Published when a workflow is about to be executed.
|
||||
/// </summary>
|
||||
public class WorkflowExecuting : WorkflowNotification
|
||||
{
|
||||
public WorkflowExecuting(WorkflowExecutionContext workflowExecutionContext) : base(workflowExecutionContext)
|
||||
{
|
||||
}
|
||||
}
|
||||
}
|
||||
46
src/core/Elsa.Core/Handlers/PersistWorkflowContext.cs
Normal file
46
src/core/Elsa.Core/Handlers/PersistWorkflowContext.cs
Normal file
|
|
@ -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<WorkflowExecuted>, INotificationHandler<ActivityExecuted>
|
||||
{
|
||||
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<string?> 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
51
src/core/Elsa.Core/Handlers/RefreshWorkflowContext.cs
Normal file
51
src/core/Elsa.Core/Handlers/RefreshWorkflowContext.cs
Normal file
|
|
@ -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<WorkflowExecuting>, INotificationHandler<ActivityExecuting>
|
||||
{
|
||||
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<object?> 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -30,7 +30,6 @@ namespace Elsa.Services
|
|||
private readonly IWorkflowSelector _workflowSelector;
|
||||
private readonly IWorkflowInstanceStore _workflowInstanceManager;
|
||||
private readonly Func<IWorkflowBuilder> _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<IWorkflowBuilder> workflowBuilderFactory,
|
||||
IWorkflowContextManager workflowContextManager,
|
||||
IMediator mediator,
|
||||
IServiceProvider serviceProvider,
|
||||
ILogger<WorkflowRunner> 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<object?> 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<string?> 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();
|
||||
|
|
|
|||
Loading…
Reference in a new issue