* Refactor workflow dispatch and instance creation process This update splits the process of dispatching a workflow into two steps: Initialization and Execution. Now, first, a new workflow instance is created and saved with the input parameters. Second, the new workflow instance is dispatched for execution. This process ensures that the size of the message dispatched does not exceed acceptable limits and enhances workflow dispatch efficiency. It also helps avoid data loss in case of premature process termination or failure in the initial stages of execution. * Change default workflow substatus to 'Pending' The code changes involve modifying the default WorkflowSubStatus from 'Executing' to 'Pending'. This minor adjustment is implemented to represent a more accurate initial state of a new workflow instance in the WorkflowManagement module in Elsa. * Add support for existing workflow instances in WorkflowInstance grain The change extends the WorkflowInstance grain to include an IWorkflowInstanceStore to support workflows that already exist. Additionally, the CreateWorkflowHostAsync method is enhanced to rebuild a workflow host if the instance already exists. * Simplify XML comments format * Refactor comments in `StartWorkflowHostParams` class Removed unnecessary comment tags in the `StartWorkflowHostParams` class for better readability and simplicity. Also, unused lines of code from the `IWorkflowRuntime` interface have been commented out to enhance clarity and cleanliness in the codebase. * Remove unnecessary comment in MassTransitWorkflowDispatcher The obsolete comment about attaching a version header to the message was removed from 'MassTransitWorkflowDispatcher.cs'. This comment was no longer relevant as the dispatcher no longer performs the action described.
131 lines
5.1 KiB
C#
131 lines
5.1 KiB
C#
using Elsa.MassTransit.Messages;
|
|
using Elsa.Workflows.Management.Contracts;
|
|
using Elsa.Workflows.Runtime.Contracts;
|
|
using Elsa.Workflows.Runtime.Options;
|
|
using Elsa.Workflows.Runtime.Parameters;
|
|
using JetBrains.Annotations;
|
|
using MassTransit;
|
|
|
|
namespace Elsa.MassTransit.Consumers;
|
|
|
|
/// <summary>
|
|
/// A consumer of various dispatch message types to asynchronously execute workflows.
|
|
/// </summary>
|
|
[UsedImplicitly]
|
|
public class DispatchWorkflowRequestConsumer :
|
|
IConsumer<DispatchWorkflowDefinition>,
|
|
IConsumer<DispatchWorkflowInstance>,
|
|
IConsumer<DispatchTriggerWorkflows>,
|
|
IConsumer<DispatchResumeWorkflows>
|
|
{
|
|
private readonly IWorkflowRuntime _workflowRuntime;
|
|
|
|
/// <summary>
|
|
/// Initializes a new instance of the <see cref="DispatchWorkflowRequestConsumer"/> class.
|
|
/// </summary>
|
|
public DispatchWorkflowRequestConsumer(IWorkflowRuntime workflowRuntime, IWorkflowInstanceManager workflowInstanceManager)
|
|
{
|
|
_workflowRuntime = workflowRuntime;
|
|
}
|
|
|
|
/// <inheritdoc />
|
|
public async Task Consume(ConsumeContext<DispatchWorkflowDefinition> context)
|
|
{
|
|
if (context.Message.IsExistingInstance)
|
|
await DispatchExistingWorkflowInstanceAsync(context.Message, context.CancellationToken);
|
|
else
|
|
await DispatchNewWorkflowInstanceAsync(context.Message, context.CancellationToken);
|
|
}
|
|
|
|
/// <inheritdoc />
|
|
public async Task Consume(ConsumeContext<DispatchWorkflowInstance> context)
|
|
{
|
|
var message = context.Message;
|
|
var cancellationToken = context.CancellationToken;
|
|
|
|
var options = new ResumeWorkflowRuntimeParams
|
|
{
|
|
CorrelationId = message.CorrelationId,
|
|
BookmarkId = message.BookmarkId,
|
|
ActivityId = message.ActivityId,
|
|
ActivityNodeId = message.ActivityNodeId,
|
|
ActivityInstanceId = message.ActivityInstanceId,
|
|
ActivityHash = message.ActivityHash,
|
|
Input = message.Input,
|
|
Properties = message.Properties,
|
|
CancellationTokens = cancellationToken
|
|
};
|
|
|
|
await _workflowRuntime.ResumeWorkflowAsync(message.InstanceId, options);
|
|
}
|
|
|
|
/// <inheritdoc />
|
|
public async Task Consume(ConsumeContext<DispatchTriggerWorkflows> context)
|
|
{
|
|
var message = context.Message;
|
|
var cancellationToken = context.CancellationToken;
|
|
var options = new TriggerWorkflowsOptions
|
|
{
|
|
CorrelationId = message.CorrelationId,
|
|
WorkflowInstanceId = message.WorkflowInstanceId,
|
|
ActivityInstanceId = message.ActivityInstanceId,
|
|
Input = message.Input,
|
|
Properties = message.Properties,
|
|
CancellationTokens = cancellationToken
|
|
};
|
|
await _workflowRuntime.TriggerWorkflowsAsync(message.ActivityTypeName, message.BookmarkPayload, options);
|
|
}
|
|
|
|
/// <inheritdoc />
|
|
public async Task Consume(ConsumeContext<DispatchResumeWorkflows> context)
|
|
{
|
|
var message = context.Message;
|
|
var cancellationToken = context.CancellationToken;
|
|
|
|
var options = new TriggerWorkflowsOptions
|
|
{
|
|
CorrelationId = message.CorrelationId,
|
|
WorkflowInstanceId = message.WorkflowInstanceId,
|
|
Input = message.Input,
|
|
Properties = message.Properties,
|
|
CancellationTokens = cancellationToken
|
|
};
|
|
|
|
await _workflowRuntime.ResumeWorkflowsAsync(message.ActivityTypeName, message.BookmarkPayload, options);
|
|
}
|
|
|
|
private async Task DispatchNewWorkflowInstanceAsync(DispatchWorkflowDefinition message, CancellationToken cancellationToken)
|
|
{
|
|
if (string.IsNullOrWhiteSpace(message.DefinitionId)) throw new ArgumentException("The definition ID is required when dispatching a workflow definition.");
|
|
if (message.VersionOptions == null) throw new ArgumentException("The version options are required when dispatching a workflow definition.");
|
|
|
|
var options = new StartWorkflowRuntimeParams
|
|
{
|
|
ParentWorkflowInstanceId = message.ParentWorkflowInstanceId,
|
|
CorrelationId = message.CorrelationId,
|
|
Input = message.Input,
|
|
Properties = message.Properties,
|
|
VersionOptions = message.VersionOptions.Value,
|
|
TriggerActivityId = message.TriggerActivityId,
|
|
InstanceId = message.InstanceId,
|
|
CancellationTokens = cancellationToken
|
|
};
|
|
|
|
await _workflowRuntime.TryStartWorkflowAsync(message.DefinitionId, options);
|
|
}
|
|
|
|
private async Task DispatchExistingWorkflowInstanceAsync(DispatchWorkflowDefinition message, CancellationToken cancellationToken)
|
|
{
|
|
if (string.IsNullOrWhiteSpace(message.InstanceId)) throw new ArgumentException("The instance ID is required when dispatching an existing workflow instance.");
|
|
|
|
var options = new StartWorkflowRuntimeParams
|
|
{
|
|
TriggerActivityId = message.TriggerActivityId,
|
|
IsExistingInstance = true,
|
|
InstanceId = message.InstanceId,
|
|
CancellationTokens = cancellationToken
|
|
};
|
|
|
|
await _workflowRuntime.StartWorkflowAsync(message.InstanceId, options);
|
|
}
|
|
} |