* Implement distributed lock for registry population Added distributed lock in 'PopulateRegistriesHostedService' to prevent concurrent registry updates. Also implemented a semaphore in 'DefaultWorkflowDefinitionStorePopulator' to control access to shared resources during add or update operations. This change helps to ensure the integrity and consistency of the workflow registries. * Refactor dependency injection for IDistributedLockProvider * Refactor option classes to parameter classes in workflow runtime This refactoring enhances the clarity of the Elsa Workflow runtime by renaming "options" classes to "parameters" classes. The name "options" misrepresented the classes' role and created confusion, as they are used to parameterize method calls rather than to configure services. The change applies to various workflow methods and tests across the project.
105 lines
3.7 KiB
C#
105 lines
3.7 KiB
C#
using Elsa.MassTransit.Messages;
|
|
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)
|
|
{
|
|
_workflowRuntime = workflowRuntime;
|
|
}
|
|
|
|
/// <inheritdoc />
|
|
public async Task Consume(ConsumeContext<DispatchWorkflowDefinition> context)
|
|
{
|
|
var message = context.Message;
|
|
var cancellationToken = context.CancellationToken;
|
|
var options = new StartWorkflowRuntimeParams
|
|
{
|
|
CorrelationId = message.CorrelationId,
|
|
Input = message.Input,
|
|
Properties = message.Properties,
|
|
VersionOptions = message.VersionOptions,
|
|
TriggerActivityId = message.TriggerActivityId,
|
|
InstanceId = message.InstanceId,
|
|
CancellationTokens = cancellationToken
|
|
};
|
|
|
|
await _workflowRuntime.TryStartWorkflowAsync(message.DefinitionId, options);
|
|
}
|
|
|
|
/// <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);
|
|
}
|
|
} |