Incremental work on default dispatchers in preparation for actor model implementation
This commit is contained in:
parent
e21154f7b8
commit
bd87f5b2df
|
|
@ -6,11 +6,14 @@ using Elsa.Activities.AzureServiceBus.Bookmarks;
|
|||
using Elsa.Activities.AzureServiceBus.Models;
|
||||
using Elsa.Activities.AzureServiceBus.Options;
|
||||
using Elsa.Bookmarks;
|
||||
using Elsa.Dispatch;
|
||||
using Elsa.DistributedLock;
|
||||
using Elsa.DistributedLocking;
|
||||
using Elsa.Models;
|
||||
using Elsa.Persistence;
|
||||
using Elsa.Persistence.Specifications;
|
||||
using Elsa.Services;
|
||||
using Elsa.Services.Models;
|
||||
using Elsa.Triggers;
|
||||
using Microsoft.Azure.ServiceBus;
|
||||
using Microsoft.Azure.ServiceBus.Core;
|
||||
|
|
@ -27,20 +30,18 @@ namespace Elsa.Activities.AzureServiceBus.Services
|
|||
private const string TenantId = default;
|
||||
|
||||
private readonly IMessageReceiver _messageReceiver;
|
||||
private readonly IServiceScopeFactory _serviceScopeFactory;
|
||||
private readonly IDistributedLockProvider _distributedLockProvider;
|
||||
private readonly ICorrelatingWorkflowDispatcher _workflowDispatcher;
|
||||
private readonly ILogger _logger;
|
||||
|
||||
public QueueWorker(
|
||||
IMessageReceiver messageReceiver,
|
||||
ICorrelatingWorkflowDispatcher workflowDispatcher,
|
||||
IServiceScopeFactory serviceScopeFactory,
|
||||
IDistributedLockProvider distributedLockProvider,
|
||||
IOptions<AzureServiceBusOptions> options,
|
||||
ILogger<QueueWorker> logger)
|
||||
{
|
||||
_messageReceiver = messageReceiver;
|
||||
_serviceScopeFactory = serviceScopeFactory;
|
||||
_distributedLockProvider = distributedLockProvider;
|
||||
_workflowDispatcher = workflowDispatcher;
|
||||
_logger = logger;
|
||||
|
||||
_messageReceiver.RegisterMessageHandler(OnMessageReceived, new MessageHandlerOptions(ExceptionReceivedHandler)
|
||||
|
|
@ -61,11 +62,8 @@ namespace Elsa.Activities.AzureServiceBus.Services
|
|||
|
||||
private async Task TriggerWorkflowsAsync(Message message, CancellationToken cancellationToken)
|
||||
{
|
||||
using var scope = _serviceScopeFactory.CreateScope();
|
||||
var workflowQueue = scope.ServiceProvider.GetRequiredService<IWorkflowQueue>();
|
||||
var queueName = _messageReceiver.Path;
|
||||
var correlationId = message.CorrelationId;
|
||||
var triggerFinder = scope.ServiceProvider.GetRequiredService<ITriggerFinder>();
|
||||
|
||||
var model = new MessageModel
|
||||
{
|
||||
|
|
@ -85,63 +83,10 @@ namespace Elsa.Activities.AzureServiceBus.Services
|
|||
ScheduledEnqueueTimeUtc = message.ScheduledEnqueueTimeUtc
|
||||
};
|
||||
|
||||
async Task TriggerNewWorkflowAsync()
|
||||
{
|
||||
var bookmark = new QueueMessageReceivedBookmark(queueName);
|
||||
var triggers = await triggerFinder.FindTriggersAsync<AzureServiceBusQueueMessageReceived>(bookmark, TenantId, cancellationToken);
|
||||
|
||||
foreach (var trigger in triggers)
|
||||
{
|
||||
var workflowBlueprint = trigger.WorkflowBlueprint;
|
||||
await workflowQueue.EnqueueWorkflowDefinition(workflowBlueprint.Id, workflowBlueprint.TenantId, trigger.ActivityId, model, correlationId, null, cancellationToken);
|
||||
}
|
||||
}
|
||||
|
||||
if (string.IsNullOrWhiteSpace(correlationId))
|
||||
{
|
||||
await TriggerNewWorkflowAsync();
|
||||
return;
|
||||
}
|
||||
|
||||
var lockKey = $"azure-service-bus:{queueName}:correlation-{correlationId}";
|
||||
var stopwatch = new Stopwatch();
|
||||
|
||||
_logger.LogDebug("Acquiring lock {LockKey}", lockKey);
|
||||
stopwatch.Start();
|
||||
|
||||
if (!await _distributedLockProvider.AcquireLockAsync(lockKey, cancellationToken))
|
||||
{
|
||||
_logger.LogDebug("Lock {LockKey} already taken", lockKey);
|
||||
return;
|
||||
}
|
||||
|
||||
try
|
||||
{
|
||||
var bookmarkFinder = scope.ServiceProvider.GetRequiredService<IBookmarkFinder>();
|
||||
var workflowInstanceStore = scope.ServiceProvider.GetRequiredService<IWorkflowInstanceStore>();
|
||||
var correlatedWorkflowInstanceCount = await workflowInstanceStore.CountAsync(new CorrelationIdSpecification<WorkflowInstance>(model.CorrelationId), cancellationToken);
|
||||
|
||||
if (correlatedWorkflowInstanceCount > 0)
|
||||
{
|
||||
// Trigger existing workflows (if blocked on this message).
|
||||
_logger.LogDebug("{WorkflowInstanceCount} existing workflows found with correlation ID '{CorrelationId}'. Resuming them", correlatedWorkflowInstanceCount, correlationId);
|
||||
var bookmark = new QueueMessageReceivedBookmark(queueName, correlationId);
|
||||
var existingWorkflows = await bookmarkFinder.FindBookmarksAsync<AzureServiceBusQueueMessageReceived>(bookmark, TenantId, cancellationToken).ToList();
|
||||
await workflowQueue.EnqueueWorkflowsAsync(existingWorkflows, model, model.CorrelationId, cancellationToken: cancellationToken);
|
||||
}
|
||||
else
|
||||
{
|
||||
// Trigger new workflow.
|
||||
_logger.LogDebug("No existing workflows found with correlation ID '{CorrelationId}'. Starting new workflow", correlationId);
|
||||
await TriggerNewWorkflowAsync();
|
||||
}
|
||||
}
|
||||
finally
|
||||
{
|
||||
await _distributedLockProvider.ReleaseLockAsync(lockKey, cancellationToken);
|
||||
stopwatch.Stop();
|
||||
_logger.LogDebug("Lock held for {ElapseTime}", stopwatch.Elapsed);
|
||||
}
|
||||
var bookmark = new QueueMessageReceivedBookmark(queueName, correlationId);
|
||||
var trigger = new QueueMessageReceivedBookmark(queueName);
|
||||
var activityType = nameof(AzureServiceBusQueueMessageReceived);
|
||||
await _workflowDispatcher.DispatchAsync(new ExecuteCorrelatedWorkflowRequest(correlationId, bookmark, trigger, activityType, model, TenantId: TenantId), cancellationToken);
|
||||
}
|
||||
|
||||
private Task ExceptionReceivedHandler(ExceptionReceivedEventArgs e)
|
||||
|
|
|
|||
|
|
@ -7,6 +7,7 @@ using Elsa.Activities.AzureServiceBus.Models;
|
|||
using Elsa.Activities.AzureServiceBus.Options;
|
||||
using Elsa.Bookmarks;
|
||||
using Elsa.DistributedLock;
|
||||
using Elsa.DistributedLocking;
|
||||
using Elsa.Models;
|
||||
using Elsa.Persistence;
|
||||
using Elsa.Persistence.Specifications;
|
||||
|
|
|
|||
|
|
@ -7,6 +7,7 @@ using Elsa.Activities.Telnyx.Webhooks.Payloads.Abstract;
|
|||
using Elsa.Activities.Telnyx.Webhooks.Payloads.Call;
|
||||
using Elsa.Activities.Telnyx.Webhooks.Services;
|
||||
using Elsa.DistributedLock;
|
||||
using Elsa.DistributedLocking;
|
||||
using Elsa.Models;
|
||||
using Elsa.Persistence;
|
||||
using Elsa.Persistence.Specifications;
|
||||
|
|
|
|||
|
|
@ -1,37 +1,20 @@
|
|||
using System.Diagnostics;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.DistributedLock;
|
||||
using Elsa.Models;
|
||||
using Elsa.Persistence;
|
||||
using Elsa.Persistence.Specifications;
|
||||
using Elsa.Persistence.Specifications.WorkflowInstances;
|
||||
using Elsa.Services;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Dispatch;
|
||||
using Quartz;
|
||||
|
||||
namespace Elsa.Activities.Temporal.Quartz.Jobs
|
||||
{
|
||||
public class RunQuartzWorkflowJob : IJob
|
||||
{
|
||||
private readonly IWorkflowRegistry _workflowRegistry;
|
||||
private readonly IWorkflowInstanceStore _workflowInstanceStore;
|
||||
private readonly IWorkflowQueue _workflowQueue;
|
||||
private readonly IDistributedLockProvider _distributedLockProvider;
|
||||
private readonly ILogger _logger;
|
||||
private readonly Stopwatch _stopwatch = new();
|
||||
private readonly IWorkflowDefinitionDispatcher _workflowDefinitionDispatcher;
|
||||
private readonly IWorkflowInstanceDispatcher _workflowInstanceDispatcher;
|
||||
|
||||
public RunQuartzWorkflowJob(
|
||||
IWorkflowRegistry workflowRegistry,
|
||||
IWorkflowInstanceStore workflowInstanceStore,
|
||||
IWorkflowQueue workflowQueue,
|
||||
IDistributedLockProvider distributedLockProvider,
|
||||
ILogger<RunQuartzWorkflowJob> logger)
|
||||
IWorkflowDefinitionDispatcher workflowDefinitionDispatcher,
|
||||
IWorkflowInstanceDispatcher workflowInstanceDispatcher)
|
||||
{
|
||||
_workflowRegistry = workflowRegistry;
|
||||
_workflowInstanceStore = workflowInstanceStore;
|
||||
_workflowQueue = workflowQueue;
|
||||
_distributedLockProvider = distributedLockProvider;
|
||||
_logger = logger;
|
||||
_workflowDefinitionDispatcher = workflowDefinitionDispatcher;
|
||||
_workflowInstanceDispatcher = workflowInstanceDispatcher;
|
||||
}
|
||||
|
||||
public async Task Execute(IJobExecutionContext context)
|
||||
|
|
@ -42,50 +25,11 @@ namespace Elsa.Activities.Temporal.Quartz.Jobs
|
|||
var tenantId = dataMap.GetString("TenantId");
|
||||
var workflowDefinitionId = dataMap.GetString("WorkflowDefinitionId")!;
|
||||
var activityId = dataMap.GetString("ActivityId")!;
|
||||
var lockKey = (workflowInstanceId, workflowDefinitionId, activityId).GetHashCode().ToString();
|
||||
|
||||
_logger.LogDebug("Acquiring lock on {WorkflowInstanceId} / {WorkflowDefinitionId} / {ActivityId}", workflowInstanceId, workflowDefinitionId, activityId);
|
||||
|
||||
if (!await _distributedLockProvider.AcquireLockAsync(lockKey, cancellationToken))
|
||||
{
|
||||
_logger.LogDebug("Failed to acquire lock on {WorkflowInstanceId} / {WorkflowDefinitionId} / {ActivityId}", workflowInstanceId, workflowDefinitionId, activityId);
|
||||
return;
|
||||
}
|
||||
|
||||
_stopwatch.Restart();
|
||||
|
||||
try
|
||||
{
|
||||
if (workflowInstanceId == null)
|
||||
{
|
||||
var workflowBlueprint = (await _workflowRegistry.GetAsync(workflowDefinitionId, tenantId, VersionOptions.Published, cancellationToken));
|
||||
|
||||
if (workflowBlueprint == null)
|
||||
{
|
||||
_logger.LogWarning("No workflow definition {WorkflowDefinitionId} found. Make sure the scheduled workflow definition is published and enabled", workflowDefinitionId);
|
||||
return;
|
||||
}
|
||||
|
||||
if (!workflowBlueprint.IsSingleton || await GetWorkflowIsAlreadyExecutingAsync(tenantId, workflowDefinitionId) == false)
|
||||
await _workflowQueue.EnqueueWorkflowDefinition(workflowDefinitionId, tenantId, activityId, null, null, null, cancellationToken);
|
||||
}
|
||||
else
|
||||
{
|
||||
await _workflowQueue.EnqueueWorkflowInstance(workflowInstanceId, activityId, null, cancellationToken);
|
||||
}
|
||||
}
|
||||
finally
|
||||
{
|
||||
_stopwatch.Stop();
|
||||
await _distributedLockProvider.ReleaseLockAsync(lockKey, cancellationToken);
|
||||
_logger.LogDebug("Held lock on {WorkflowInstanceId} / {WorkflowDefinitionId} / {ActivityId} for {LockTime}", workflowInstanceId, workflowDefinitionId, activityId, _stopwatch.Elapsed);
|
||||
}
|
||||
}
|
||||
|
||||
private async Task<bool> GetWorkflowIsAlreadyExecutingAsync(string? tenantId, string workflowDefinitionId)
|
||||
{
|
||||
var specification = new TenantSpecification<WorkflowInstance>(tenantId).WithWorkflowDefinition(workflowDefinitionId).And(new WorkflowIsAlreadyExecutingSpecification());
|
||||
return await _workflowInstanceStore.FindAsync(specification) != null;
|
||||
if (workflowInstanceId == null)
|
||||
await _workflowDefinitionDispatcher.DispatchAsync(new ExecuteWorkflowDefinitionRequest(workflowDefinitionId, activityId, TenantId: tenantId), cancellationToken);
|
||||
else
|
||||
await _workflowInstanceDispatcher.DispatchAsync(new ExecuteWorkflowInstanceRequest(workflowInstanceId, activityId), cancellationToken);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,14 @@
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Elsa.Dispatch
|
||||
{
|
||||
/// <summary>
|
||||
/// The correlating dispatcher is responsible for finding workflows correlated by the specified correlation ID.
|
||||
/// If no correlated workflows are found, a new one is started.
|
||||
/// </summary>
|
||||
public interface ICorrelatingWorkflowDispatcher
|
||||
{
|
||||
Task DispatchAsync(ExecuteCorrelatedWorkflowRequest request, CancellationToken? cancellationToken = default);
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,13 @@
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Elsa.Dispatch
|
||||
{
|
||||
/// <summary>
|
||||
/// Dispatches requests for executing workflow definitions.
|
||||
/// </summary>
|
||||
public interface IWorkflowDefinitionDispatcher
|
||||
{
|
||||
Task DispatchAsync(ExecuteWorkflowDefinitionRequest request, CancellationToken? cancellationToken = default);
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,13 @@
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Elsa.Dispatch
|
||||
{
|
||||
/// <summary>
|
||||
/// Dispatches requests for executing workflow instances.
|
||||
/// </summary>
|
||||
public interface IWorkflowInstanceDispatcher
|
||||
{
|
||||
Task DispatchAsync(ExecuteWorkflowInstanceRequest request, CancellationToken? cancellationToken = default);
|
||||
}
|
||||
}
|
||||
9
src/core/Elsa.Abstractions/Dispatch/Models.cs
Normal file
9
src/core/Elsa.Abstractions/Dispatch/Models.cs
Normal file
|
|
@ -0,0 +1,9 @@
|
|||
using Elsa.Bookmarks;
|
||||
|
||||
namespace Elsa.Dispatch
|
||||
{
|
||||
public record ExecuteCorrelatedWorkflowRequest(string CorrelationId, IBookmark Bookmark, IBookmark Trigger, string ActivityType, object? Input = default, string? ContextId = default, string? TenantId = default);
|
||||
public record ExecuteWorkflowDefinitionRequest(string WorkflowDefinitionId, string? ActivityId = default, object? Input = default, string? CorrelationId = default, string? ContextId = default, string? TenantId = default);
|
||||
public record ExecuteWorkflowInstanceRequest(string WorkflowInstanceId, string ActivityId, object? Input = default);
|
||||
|
||||
}
|
||||
|
|
@ -1,7 +1,7 @@
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
|
||||
namespace Elsa.DistributedLock
|
||||
namespace Elsa.DistributedLocking
|
||||
{
|
||||
/// <summary>
|
||||
/// Provides functionality to acquire and release locks which are distributed across all instances in a web farm.
|
||||
|
|
@ -1,20 +1,20 @@
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Bookmarks;
|
||||
using Elsa.Services;
|
||||
|
||||
namespace Elsa
|
||||
{
|
||||
public static class WorkflowQueueExtensions
|
||||
{
|
||||
public static Task EnqueueWorkflowsAsync<T>(
|
||||
this IWorkflowQueue workflowQueue,
|
||||
IBookmark bookmark,
|
||||
string? tenantId,
|
||||
object? input = default,
|
||||
string? correlationId = default,
|
||||
string? contextId = default,
|
||||
CancellationToken cancellationToken = default) where T : IActivity =>
|
||||
workflowQueue.EnqueueWorkflowsAsync(typeof(T).Name, bookmark, tenantId, input, correlationId, contextId, cancellationToken);
|
||||
}
|
||||
}
|
||||
// using System.Threading;
|
||||
// using System.Threading.Tasks;
|
||||
// using Elsa.Bookmarks;
|
||||
// using Elsa.Services;
|
||||
//
|
||||
// namespace Elsa
|
||||
// {
|
||||
// public static class WorkflowQueueExtensions
|
||||
// {
|
||||
// public static Task EnqueueWorkflowsAsync<T>(
|
||||
// this IWorkflowQueue workflowQueue,
|
||||
// IBookmark bookmark,
|
||||
// string? tenantId,
|
||||
// object? input = default,
|
||||
// string? correlationId = default,
|
||||
// string? contextId = default,
|
||||
// CancellationToken cancellationToken = default) where T : IActivity =>
|
||||
// workflowQueue.EnqueueWorkflowsAsync(typeof(T).Name, bookmark, tenantId, input, correlationId, contextId, cancellationToken);
|
||||
// }
|
||||
// }
|
||||
|
|
@ -1,5 +1,7 @@
|
|||
using System;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.DistributedLock;
|
||||
using Elsa.DistributedLocking;
|
||||
using Elsa.Messages;
|
||||
using Elsa.Models;
|
||||
using Elsa.Persistence;
|
||||
|
|
@ -12,6 +14,7 @@ using Rebus.Handlers;
|
|||
|
||||
namespace Elsa.Consumers
|
||||
{
|
||||
[Obsolete]
|
||||
public class RunWorkflowDefinitionConsumer : IHandleMessages<RunWorkflowDefinition>
|
||||
{
|
||||
private readonly IWorkflowRunner _workflowRunner;
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ using System.Diagnostics;
|
|||
using System.Linq;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.DistributedLock;
|
||||
using Elsa.DistributedLocking;
|
||||
using Elsa.Messages;
|
||||
using Elsa.Models;
|
||||
using Elsa.Persistence;
|
||||
|
|
|
|||
|
|
@ -0,0 +1,106 @@
|
|||
using System.Collections.Generic;
|
||||
using System.Diagnostics;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Bookmarks;
|
||||
using Elsa.DistributedLocking;
|
||||
using Elsa.Models;
|
||||
using Elsa.Persistence;
|
||||
using Elsa.Persistence.Specifications;
|
||||
using Elsa.Services;
|
||||
using Elsa.Triggers;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using Open.Linq.AsyncExtensions;
|
||||
using Rebus.Handlers;
|
||||
|
||||
namespace Elsa.Dispatch.Consumers
|
||||
{
|
||||
public class ExecuteCorrelatedWorkflowRequestConsumer : IHandleMessages<ExecuteCorrelatedWorkflowRequest>
|
||||
{
|
||||
private readonly IWorkflowInstanceStore _workflowInstanceStore;
|
||||
private readonly IDistributedLockProvider _distributedLockProvider;
|
||||
private readonly IBookmarkFinder _bookmarkFinder;
|
||||
private readonly ITriggerFinder _triggerFinder;
|
||||
private readonly ICommandSender _commandSender;
|
||||
private readonly ILogger _logger;
|
||||
private readonly Stopwatch _stopwatch = new();
|
||||
|
||||
public ExecuteCorrelatedWorkflowRequestConsumer(
|
||||
IWorkflowInstanceStore workflowInstanceStore,
|
||||
IDistributedLockProvider distributedLockProvider,
|
||||
IBookmarkFinder bookmarkFinder,
|
||||
ITriggerFinder triggerFinder,
|
||||
ICommandSender commandSender,
|
||||
ILogger<ExecuteCorrelatedWorkflowRequestConsumer> logger)
|
||||
{
|
||||
_workflowInstanceStore = workflowInstanceStore;
|
||||
_distributedLockProvider = distributedLockProvider;
|
||||
_bookmarkFinder = bookmarkFinder;
|
||||
_triggerFinder = triggerFinder;
|
||||
_commandSender = commandSender;
|
||||
_logger = logger;
|
||||
}
|
||||
|
||||
public async Task Handle(ExecuteCorrelatedWorkflowRequest message)
|
||||
{
|
||||
var correlationId = message.CorrelationId;
|
||||
var lockKey = $"correlated-workflow-request:correlation-{correlationId}";
|
||||
|
||||
_logger.LogDebug("Acquiring lock {LockKey}", lockKey);
|
||||
_stopwatch.Restart();
|
||||
|
||||
if (!await _distributedLockProvider.AcquireLockAsync(lockKey))
|
||||
{
|
||||
_logger.LogDebug("Lock {LockKey} already taken", lockKey);
|
||||
await _commandSender.SendAsync(message);
|
||||
return;
|
||||
}
|
||||
|
||||
try
|
||||
{
|
||||
var correlatedWorkflowInstanceCount = await _workflowInstanceStore.CountAsync(new CorrelationIdSpecification<WorkflowInstance>(correlationId));
|
||||
|
||||
if (correlatedWorkflowInstanceCount > 0)
|
||||
{
|
||||
_logger.LogDebug("{WorkflowInstanceCount} existing workflows found with correlation ID '{CorrelationId}' will be queued for execution", correlatedWorkflowInstanceCount, correlationId);
|
||||
var existingWorkflows = await _bookmarkFinder.FindBookmarksAsync(message.ActivityType, message.Bookmark, message.TenantId).ToList();
|
||||
await EnqueueWorkflowsAsync(existingWorkflows, message.Input);
|
||||
}
|
||||
else
|
||||
{
|
||||
// Trigger new workflow.
|
||||
_logger.LogDebug("No existing workflows found with correlation ID '{CorrelationId}'. Starting new workflow", correlationId);
|
||||
await TriggerNewWorkflowAsync(message);
|
||||
}
|
||||
}
|
||||
finally
|
||||
{
|
||||
await _distributedLockProvider.ReleaseLockAsync(lockKey);
|
||||
_stopwatch.Stop();
|
||||
_logger.LogDebug("Lock held for {ElapseTime}", _stopwatch.Elapsed);
|
||||
}
|
||||
}
|
||||
|
||||
async Task TriggerNewWorkflowAsync(ExecuteCorrelatedWorkflowRequest message)
|
||||
{
|
||||
var filter = message.Trigger;
|
||||
var triggers = await _triggerFinder.FindTriggersAsync(message.ActivityType, filter, message.TenantId);
|
||||
|
||||
foreach (var trigger in triggers)
|
||||
{
|
||||
var workflowBlueprint = trigger.WorkflowBlueprint;
|
||||
await EnqueueWorkflowDefinition(workflowBlueprint.Id, workflowBlueprint.TenantId, trigger.ActivityId, message.Input, message.CorrelationId, message.ContextId);
|
||||
}
|
||||
}
|
||||
|
||||
private async Task EnqueueWorkflowsAsync(IEnumerable<BookmarkFinderResult> results, object? input)
|
||||
{
|
||||
foreach (var result in results)
|
||||
await EnqueueWorkflowInstance(result.WorkflowInstanceId, result.ActivityId, input);
|
||||
}
|
||||
|
||||
public async Task EnqueueWorkflowInstance(string workflowInstanceId, string activityId, object? input) => await _commandSender.SendAsync(new ExecuteWorkflowInstanceRequest(workflowInstanceId, activityId, input));
|
||||
|
||||
public async Task EnqueueWorkflowDefinition(string workflowDefinitionId, string? tenantId, string activityId, object? input, string? correlationId, string? contextId) =>
|
||||
await _commandSender.SendAsync(new ExecuteWorkflowDefinitionRequest(workflowDefinitionId, activityId, input, correlationId, contextId, tenantId));
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,62 @@
|
|||
using System.Threading.Tasks;
|
||||
using Elsa.Models;
|
||||
using Elsa.Persistence;
|
||||
using Elsa.Persistence.Specifications;
|
||||
using Elsa.Persistence.Specifications.WorkflowInstances;
|
||||
using Elsa.Services;
|
||||
using Elsa.Services.Models;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using Rebus.Handlers;
|
||||
|
||||
namespace Elsa.Dispatch.Consumers
|
||||
{
|
||||
public class ExecuteWorkflowDefinitionRequestConsumer : IHandleMessages<ExecuteWorkflowDefinitionRequest>
|
||||
{
|
||||
private readonly IWorkflowRunner _workflowRunner;
|
||||
private readonly IWorkflowRegistry _workflowRegistry;
|
||||
private readonly IWorkflowInstanceStore _workflowInstanceStore;
|
||||
private readonly ILogger _logger;
|
||||
|
||||
public ExecuteWorkflowDefinitionRequestConsumer(
|
||||
IWorkflowRunner workflowRunner,
|
||||
IWorkflowRegistry workflowRegistry,
|
||||
IWorkflowInstanceStore workflowInstanceStore,
|
||||
ILogger<ExecuteWorkflowDefinitionRequestConsumer> logger)
|
||||
{
|
||||
_workflowRunner = workflowRunner;
|
||||
_workflowRegistry = workflowRegistry;
|
||||
_workflowInstanceStore = workflowInstanceStore;
|
||||
_logger = logger;
|
||||
}
|
||||
|
||||
public async Task Handle(ExecuteWorkflowDefinitionRequest message)
|
||||
{
|
||||
var workflowDefinitionId = message.WorkflowDefinitionId;
|
||||
var tenantId = message.TenantId;
|
||||
var workflowBlueprint = await _workflowRegistry.GetAsync(workflowDefinitionId, tenantId, VersionOptions.Published);
|
||||
|
||||
if (!ValidatePreconditions(workflowDefinitionId, workflowBlueprint))
|
||||
return;
|
||||
|
||||
if (!workflowBlueprint!.IsSingleton || await GetWorkflowIsAlreadyExecutingAsync(tenantId, workflowDefinitionId) == false)
|
||||
await _workflowRunner.RunWorkflowAsync(workflowBlueprint, message.ActivityId, message.Input, message.CorrelationId, message.ContextId);
|
||||
}
|
||||
|
||||
private bool ValidatePreconditions(string? workflowDefinitionId, IWorkflowBlueprint? workflowBlueprint)
|
||||
{
|
||||
if (workflowBlueprint == null)
|
||||
{
|
||||
_logger.LogWarning("No workflow definition {WorkflowDefinitionId} found. Make sure the scheduled workflow definition is published and enabled", workflowDefinitionId);
|
||||
return false;
|
||||
}
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
private async Task<bool> GetWorkflowIsAlreadyExecutingAsync(string? tenantId, string workflowDefinitionId)
|
||||
{
|
||||
var specification = new TenantSpecification<WorkflowInstance>(tenantId).WithWorkflowDefinition(workflowDefinitionId).And(new WorkflowIsAlreadyExecutingSpecification());
|
||||
return await _workflowInstanceStore.FindAsync(specification) != null;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,100 @@
|
|||
using System.Diagnostics;
|
||||
using System.Linq;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.DistributedLocking;
|
||||
using Elsa.Models;
|
||||
using Elsa.Persistence;
|
||||
using Elsa.Services;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using Rebus.Handlers;
|
||||
|
||||
namespace Elsa.Dispatch.Consumers
|
||||
{
|
||||
public class ExecuteWorkflowRequestConsumer : IHandleMessages<ExecuteWorkflowInstanceRequest>
|
||||
{
|
||||
private readonly IWorkflowRunner _workflowRunner;
|
||||
private readonly IWorkflowInstanceStore _workflowInstanceStore;
|
||||
private readonly IDistributedLockProvider _distributedLockProvider;
|
||||
private readonly ICommandSender _commandSender;
|
||||
private readonly ILogger _logger;
|
||||
private readonly Stopwatch _stopwatch = new();
|
||||
|
||||
public ExecuteWorkflowRequestConsumer(
|
||||
IWorkflowRunner workflowRunner,
|
||||
IWorkflowInstanceStore workflowInstanceStore,
|
||||
IDistributedLockProvider distributedLockProvider,
|
||||
ICommandSender commandSender,
|
||||
ILogger<ExecuteWorkflowRequestConsumer> logger)
|
||||
{
|
||||
_workflowRunner = workflowRunner;
|
||||
_workflowInstanceStore = workflowInstanceStore;
|
||||
_distributedLockProvider = distributedLockProvider;
|
||||
_commandSender = commandSender;
|
||||
_logger = logger;
|
||||
}
|
||||
|
||||
public async Task Handle(ExecuteWorkflowInstanceRequest message)
|
||||
{
|
||||
var workflowInstanceId = message.WorkflowInstanceId;
|
||||
var lockKey = workflowInstanceId;
|
||||
|
||||
_logger.LogDebug("Acquiring lock on workflow instance {WorkflowInstanceId}", workflowInstanceId);
|
||||
_stopwatch.Restart();
|
||||
|
||||
if (!await _distributedLockProvider.AcquireLockAsync(lockKey))
|
||||
{
|
||||
_logger.LogDebug("Failed to acquire lock on workflow instance {WorkflowInstanceId}. Re-queueing message", workflowInstanceId);
|
||||
await _commandSender.SendAsync(message);
|
||||
return;
|
||||
}
|
||||
|
||||
try
|
||||
{
|
||||
var workflowInstance = await _workflowInstanceStore.FindByIdAsync(message.WorkflowInstanceId);
|
||||
|
||||
if (!ValidatePreconditions(workflowInstanceId, workflowInstance, message.ActivityId))
|
||||
return;
|
||||
|
||||
await _workflowRunner.RunWorkflowAsync(
|
||||
workflowInstance!,
|
||||
message.ActivityId,
|
||||
message.Input);
|
||||
}
|
||||
finally
|
||||
{
|
||||
await _distributedLockProvider.ReleaseLockAsync(lockKey);
|
||||
_stopwatch.Stop();
|
||||
_logger.LogDebug("Held lock on workflow instance {WorkflowInstanceId} for {ElapsedTime}", workflowInstanceId, _stopwatch.Elapsed);
|
||||
}
|
||||
}
|
||||
|
||||
private bool ValidatePreconditions(string? workflowInstanceId, WorkflowInstance? workflowInstance, string? activityId)
|
||||
{
|
||||
if (workflowInstance == null)
|
||||
{
|
||||
_logger.LogWarning("Could not run workflow instance with ID {WorkflowInstanceId} because it does not exist", workflowInstanceId);
|
||||
return false;
|
||||
}
|
||||
|
||||
if (workflowInstance.WorkflowStatus != WorkflowStatus.Suspended && workflowInstance.WorkflowStatus != WorkflowStatus.Running)
|
||||
{
|
||||
_logger.LogWarning("Could not run workflow instance with ID {WorkflowInstanceId} because it has a status other than Suspended or Running. Its actual status is {WorkflowStatus}", workflowInstanceId, workflowInstance.WorkflowStatus);
|
||||
return false;
|
||||
}
|
||||
|
||||
if (activityId != null)
|
||||
{
|
||||
var activityIsBlocking = workflowInstance.BlockingActivities.Any(x => x.ActivityId == activityId);
|
||||
var activityIsScheduled = workflowInstance.ScheduledActivities.Any(x => x.ActivityId == activityId) || workflowInstance.CurrentActivity?.ActivityId == activityId;
|
||||
|
||||
if (!activityIsBlocking && !activityIsScheduled)
|
||||
{
|
||||
_logger.LogWarning("Did not run workflow {WorkflowInstanceId} for activity {ActivityId} because the workflow is not blocked on that activity nor is that activity scheduled for execution", workflowInstanceId, activityId);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
return true;
|
||||
}
|
||||
}
|
||||
}
|
||||
18
src/core/Elsa.Core/Dispatch/QueuingWorkflowDispatcher.cs
Normal file
18
src/core/Elsa.Core/Dispatch/QueuingWorkflowDispatcher.cs
Normal file
|
|
@ -0,0 +1,18 @@
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.Services;
|
||||
|
||||
namespace Elsa.Dispatch
|
||||
{
|
||||
/// <summary>
|
||||
/// The default strategy that process workflow execution requests by sending them to a queue.
|
||||
/// </summary>
|
||||
public class QueuingWorkflowDispatcher : IWorkflowDefinitionDispatcher, IWorkflowInstanceDispatcher, ICorrelatingWorkflowDispatcher
|
||||
{
|
||||
private readonly ICommandSender _commandSender;
|
||||
public QueuingWorkflowDispatcher(ICommandSender commandSender) => _commandSender = commandSender;
|
||||
public async Task DispatchAsync(ExecuteWorkflowInstanceRequest request, CancellationToken? cancellationToken = default) => await _commandSender.SendAsync(request);
|
||||
public async Task DispatchAsync(ExecuteCorrelatedWorkflowRequest request, CancellationToken? cancellationToken = default) => await _commandSender.SendAsync(request);
|
||||
public async Task DispatchAsync(ExecuteWorkflowDefinitionRequest request, CancellationToken? cancellationToken = default) => await _commandSender.SendAsync(request);
|
||||
}
|
||||
}
|
||||
|
|
@ -1,6 +1,7 @@
|
|||
using System.Collections.Generic;
|
||||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.DistributedLocking;
|
||||
|
||||
namespace Elsa.DistributedLock
|
||||
{
|
||||
|
|
|
|||
|
|
@ -6,6 +6,7 @@ using AutoMapper;
|
|||
using Elsa.Builders;
|
||||
using Elsa.Caching;
|
||||
using Elsa.DistributedLock;
|
||||
using Elsa.DistributedLocking;
|
||||
using Elsa.Models;
|
||||
using Elsa.Persistence;
|
||||
using Elsa.Persistence.InMemory;
|
||||
|
|
|
|||
|
|
@ -2,6 +2,7 @@ using System.Linq;
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.DistributedLock;
|
||||
using Elsa.DistributedLocking;
|
||||
using Elsa.Models;
|
||||
using Elsa.Persistence;
|
||||
using Elsa.Persistence.Specifications.WorkflowInstances;
|
||||
|
|
|
|||
|
|
@ -5,6 +5,7 @@ using System.Linq;
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.DistributedLock;
|
||||
using Elsa.DistributedLocking;
|
||||
using Microsoft.Azure.Storage;
|
||||
using Microsoft.Azure.Storage.Blob;
|
||||
using Microsoft.Azure.Storage.RetryPolicies;
|
||||
|
|
|
|||
|
|
@ -3,6 +3,7 @@ using System.Collections.Generic;
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.DistributedLock;
|
||||
using Elsa.DistributedLocking;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using RedLockNet;
|
||||
|
||||
|
|
|
|||
|
|
@ -5,6 +5,7 @@ using System.Data.SqlClient;
|
|||
using System.Threading;
|
||||
using System.Threading.Tasks;
|
||||
using Elsa.DistributedLock;
|
||||
using Elsa.DistributedLocking;
|
||||
using Microsoft.Extensions.Logging;
|
||||
|
||||
namespace Elsa
|
||||
|
|
|
|||
Loading…
Reference in a new issue