Add workflow cancellation notifications (#6075)
Introduced WorkflowCancelled and WorkflowCancelling notification records. Changed cancellation service to return a boolean instead of an integer. Integrated mediator notifications in the workflow canceler service.
This commit is contained in:
parent
83860db22c
commit
6093eea01b
|
|
@ -16,7 +16,6 @@ public class DispatchCancelWorkflowsRequestConsumer(IWorkflowRuntime workflowRun
|
|||
{
|
||||
var cancellationToken = context.CancellationToken;
|
||||
var request = context.Message;
|
||||
|
||||
var client = await workflowRuntime.CreateClientAsync(request.WorkflowInstanceId, cancellationToken);
|
||||
await client.CancelAsync(cancellationToken);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -11,7 +11,7 @@ public interface IWorkflowCancellationService
|
|||
/// Cancels a workflow instance.
|
||||
/// </summary>
|
||||
/// <remarks>Also cancels all children</remarks>
|
||||
Task<int> CancelWorkflowAsync(string workflowInstanceId, CancellationToken cancellationToken = default);
|
||||
Task<bool> CancelWorkflowAsync(string workflowInstanceId, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// Cancels workflow executions with the specified workflow instance ID.
|
||||
|
|
|
|||
|
|
@ -0,0 +1,5 @@
|
|||
using Elsa.Mediator.Contracts;
|
||||
|
||||
namespace Elsa.Workflows.Runtime.Notifications;
|
||||
|
||||
public record WorkflowCancelled(string WorkflowInstanceId) : INotification;
|
||||
|
|
@ -0,0 +1,5 @@
|
|||
using Elsa.Mediator.Contracts;
|
||||
|
||||
namespace Elsa.Workflows.Runtime.Notifications;
|
||||
|
||||
public record WorkflowCancelling(string WorkflowInstanceId) : INotification;
|
||||
|
|
@ -1,39 +1,36 @@
|
|||
using Elsa.Mediator.Contracts;
|
||||
using Elsa.Workflows.Models;
|
||||
using Elsa.Workflows.Pipelines.WorkflowExecution;
|
||||
using Elsa.Workflows.Runtime.Middleware.Workflows;
|
||||
using Elsa.Workflows.Runtime.Notifications;
|
||||
using Elsa.Workflows.State;
|
||||
|
||||
namespace Elsa.Workflows.Runtime;
|
||||
|
||||
/// <inheritdoc />
|
||||
public class WorkflowCanceler(IWorkflowExecutionPipeline workflowExecutionPipeline, IWorkflowStateExtractor workflowStateExtractor, IServiceProvider serviceProvider) : IWorkflowCanceler
|
||||
public class WorkflowCanceler(
|
||||
IWorkflowExecutionPipeline workflowExecutionPipeline,
|
||||
IWorkflowStateExtractor workflowStateExtractor,
|
||||
IMediator mediator,
|
||||
IServiceProvider serviceProvider) : IWorkflowCanceler
|
||||
{
|
||||
/// <inheritdoc />
|
||||
public async Task<WorkflowState> CancelWorkflowAsync(WorkflowGraph workflowGraph, WorkflowState workflowState, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync(serviceProvider, workflowGraph, workflowState, cancellationToken: cancellationToken);
|
||||
|
||||
// Alter the workflow execution context to cancel the workflow.
|
||||
await CancelWorkflowAsync(workflowExecutionContext, cancellationToken);
|
||||
|
||||
// Map the workflow execution context back to a workflow state.
|
||||
return workflowStateExtractor.Extract(workflowExecutionContext);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task CancelWorkflowAsync(WorkflowExecutionContext workflowExecutionContext, CancellationToken cancellationToken = default)
|
||||
{
|
||||
// Build a new workflow execution pipeline.
|
||||
await mediator.SendAsync(new WorkflowCancelling(workflowExecutionContext.Id), cancellationToken);
|
||||
var pipelineBuilder = new WorkflowExecutionPipelineBuilder(serviceProvider);
|
||||
workflowExecutionPipeline.ConfigurePipelineBuilder(pipelineBuilder);
|
||||
|
||||
// Replace the terminal DefaultActivitySchedulerMiddleware with the CancelWorkflowMiddleware terminal.
|
||||
pipelineBuilder.ReplaceTerminal<CancelWorkflowMiddleware>();
|
||||
|
||||
// Build modified pipeline.
|
||||
var pipeline = pipelineBuilder.Build();
|
||||
|
||||
// Execute the pipeline.
|
||||
await pipeline(workflowExecutionContext);
|
||||
await mediator.SendAsync(new WorkflowCancelled(workflowExecutionContext.Id), cancellationToken);
|
||||
}
|
||||
}
|
||||
|
|
@ -15,7 +15,7 @@ public class WorkflowCancellationService(
|
|||
: IWorkflowCancellationService
|
||||
{
|
||||
/// <inheritdoc />
|
||||
public async Task<int> CancelWorkflowAsync(string workflowInstanceId, CancellationToken cancellationToken = default)
|
||||
public async Task<bool> CancelWorkflowAsync(string workflowInstanceId, CancellationToken cancellationToken = default)
|
||||
{
|
||||
var filter = new WorkflowInstanceFilter
|
||||
{
|
||||
|
|
@ -23,12 +23,11 @@ public class WorkflowCancellationService(
|
|||
};
|
||||
var instance = await workflowInstanceStore.FindAsync(filter, cancellationToken);
|
||||
|
||||
return instance == null
|
||||
? 0
|
||||
: await CancelWorkflows(new List<WorkflowInstance>
|
||||
{
|
||||
instance
|
||||
}, cancellationToken);
|
||||
if(instance == null)
|
||||
return false;
|
||||
|
||||
await CancelWorkflows([instance], cancellationToken);
|
||||
return true;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
|
|
|
|||
Loading…
Reference in a new issue