Fix deadlock + multiple trigger support (#3720)

* Lock on specific hash for new workflows and split between trigger new /& existing workflows

* Update flowchart to take into account the triggered activity to start

* Use dispatcher for resuming parent workflows

* Update logging middleware (and remove it by default)
This commit is contained in:
Sipke Schoorstra 2023-02-20 12:05:57 +01:00 committed by GitHub
parent 0a831ff65b
commit d2d7f23d73
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
11 changed files with 126 additions and 99 deletions

View file

@ -7,15 +7,22 @@ using Elsa.Workflows.Runtime.Services;
namespace Elsa.Jobs.Activities.Handlers;
/// <summary>
/// A handler that resumes workflows waiting for a given job to complete.
/// </summary>
public class JobExecutedHandler : INotificationHandler<JobExecuted>
{
private readonly IWorkflowRuntime _workflowRuntime;
/// <summary>
/// Constructor.
/// </summary>
public JobExecutedHandler(IWorkflowRuntime workflowRuntime)
{
_workflowRuntime = workflowRuntime;
}
/// <inheritdoc />
public async Task HandleAsync(JobExecuted notification, CancellationToken cancellationToken)
{
var job = notification.Job;
@ -33,7 +40,7 @@ public class JobExecutedHandler : INotificationHandler<JobExecuted>
var jobType = notification.Job.GetType();
var payload = new EnqueuedJobPayload(job.Id);
var jobTypeName = JobTypeNameHelper.GenerateTypeName(jobType);
await _workflowRuntime.ResumeWorkflowsAsync(jobTypeName, payload, new ResumeWorkflowRuntimeOptions(), cancellationToken);
await _workflowRuntime.ResumeWorkflowsAsync(jobTypeName, payload, new TriggerWorkflowsRuntimeOptions(), cancellationToken);
}
}
}

View file

@ -52,7 +52,7 @@ public class DispatchWorkflowRequestConsumer :
public async Task Consume(ConsumeContext<DispatchResumeWorkflows> context)
{
var message = context.Message;
var options = new ResumeWorkflowRuntimeOptions(CorrelationId: message.CorrelationId, Input: message.Input);
var options = new TriggerWorkflowsRuntimeOptions(CorrelationId: message.CorrelationId, Input: message.Input);
await _workflowRuntime.ResumeWorkflowsAsync(message.ActivityTypeName, message.BookmarkPayload, options, context.CancellationToken);
}
}

View file

@ -89,6 +89,31 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
return new WorkflowExecutionResult(workflowInstanceId, bookmarks);
}
/// <inheritdoc />
public async Task<ICollection<WorkflowExecutionResult>> StartWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default)
{
var hash = _hasher.Hash(activityTypeName, bookmarkPayload);
var triggers = await _triggerStore.FindAsync(hash, cancellationToken);
var results = new List<WorkflowExecutionResult>();
foreach (var trigger in triggers)
{
var definitionId = trigger.WorkflowDefinitionId;
var startOptions = new StartWorkflowRuntimeOptions(options.CorrelationId, options.Input, VersionOptions.Published, trigger.ActivityId);
var canStartResult = await CanStartWorkflowAsync(definitionId, startOptions, cancellationToken);
// If we can't start the workflow, don't try it.
if(!canStartResult.CanStart)
continue;
var startResult = await StartWorkflowAsync(definitionId, startOptions, cancellationToken);
results.Add(new WorkflowExecutionResult(startResult.InstanceId, startResult.Bookmarks));
}
return results;
}
/// <inheritdoc />
public async Task<ResumeWorkflowResult> ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
{
@ -109,7 +134,7 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
}
/// <inheritdoc />
public async Task<ICollection<WorkflowExecutionResult>> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
public async Task<ICollection<WorkflowExecutionResult>> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default)
{
var hash = _hasher.Hash(activityTypeName, bookmarkPayload);
var client = _cluster.GetNamedBookmarkGrain(hash);
@ -122,58 +147,17 @@ public class ProtoActorWorkflowRuntime : IWorkflowRuntime
var bookmarksResponse = await client.Resolve(request, cancellationToken);
var bookmarks = bookmarksResponse!.Bookmarks;
return await ResumeWorkflowsAsync(bookmarks, options, cancellationToken);
return await ResumeWorkflowsAsync(bookmarks, new ResumeWorkflowRuntimeOptions(options.CorrelationId, Input: options.Input), cancellationToken);
}
/// <inheritdoc />
public async Task<TriggerWorkflowsResult> TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default)
{
var triggeredWorkflows = new List<WorkflowExecutionResult>();
var hash = _hasher.Hash(activityTypeName, bookmarkPayload);
var startedWorkflows = await StartWorkflowsAsync(activityTypeName, bookmarkPayload, options, cancellationToken);
var resumedWorkflows = await ResumeWorkflowsAsync(activityTypeName, bookmarkPayload, options, cancellationToken);
var results = startedWorkflows.Concat(resumedWorkflows).ToList();
// Start new workflows.
var triggers = await _triggerStore.FindAsync(hash, cancellationToken);
foreach (var trigger in triggers)
{
var definitionId = trigger.WorkflowDefinitionId;
var startOptions = new StartWorkflowRuntimeOptions(options.CorrelationId, options.Input, VersionOptions.Published, trigger.ActivityId);
var canStartResult = await CanStartWorkflowAsync(definitionId, startOptions, cancellationToken);
// If we can't start the workflow, don't try it.
if(!canStartResult.CanStart)
continue;
var startResult = await StartWorkflowAsync(definitionId, startOptions, cancellationToken);
triggeredWorkflows.Add(new WorkflowExecutionResult(startResult.InstanceId, startResult.Bookmarks));
}
// Resume existing workflow instances.
var client = _cluster.GetNamedBookmarkGrain(hash);
var request = new ResolveBookmarksRequest
{
ActivityTypeName = activityTypeName,
CorrelationId = options.CorrelationId.EmptyIfNull()
};
var bookmarksResponse = await client.Resolve(request, cancellationToken);
var bookmarks = bookmarksResponse!.Bookmarks;
foreach (var bookmark in bookmarks)
{
var workflowInstanceId = bookmark.WorkflowInstanceId;
var resumeResult = await ResumeWorkflowAsync(
workflowInstanceId,
new ResumeWorkflowRuntimeOptions(options.CorrelationId, bookmark.BookmarkId, null, options.Input),
cancellationToken);
triggeredWorkflows.Add(new WorkflowExecutionResult(workflowInstanceId, resumeResult.Bookmarks));
}
return new TriggerWorkflowsResult(triggeredWorkflows);
return new TriggerWorkflowsResult(results);
}
/// <inheritdoc />

View file

@ -36,7 +36,7 @@ internal class TriggerWebhookDrivenActivities : INotificationHandler<TelnyxWebho
var bookmarkPayload = new WebhookEventBookmarkPayload(eventType);
foreach (var activityDescriptor in activityDescriptors)
await _workflowRuntime.ResumeWorkflowsAsync(activityDescriptor.TypeName, bookmarkPayload, new ResumeWorkflowRuntimeOptions(correlationId, Input: input), cancellationToken);
await _workflowRuntime.ResumeWorkflowsAsync(activityDescriptor.TypeName, bookmarkPayload, new TriggerWorkflowsRuntimeOptions(correlationId, Input: input), cancellationToken);
}
private IEnumerable<ActivityDescriptor> FindActivityDescriptors(string eventType) =>

View file

@ -27,17 +27,25 @@ public class Flowchart : Container
}
[Port] [Browsable(false)] public IActivity? Start { get; set; }
/// <summary>
/// A list of connections between activities.
/// </summary>
public ICollection<Connection> Connections { get; set; } = new List<Connection>();
/// <inheritdoc />
protected override async ValueTask ScheduleChildrenAsync(ActivityExecutionContext context)
{
if (Start == null!)
var triggerActivityId = context.WorkflowExecutionContext.TriggerActivityId;
var triggerActivity = triggerActivityId != null ? Activities.FirstOrDefault(x => x.Id == triggerActivityId) : default;
var startActivity = triggerActivity ?? Start;
if (startActivity == null!)
{
await context.CompleteActivityAsync();
return;
}
await context.ScheduleActivityAsync(Start);
await context.ScheduleActivityAsync(startActivity);
}
private async ValueTask OnDescendantCompletedAsync(ActivityCompleted signal, SignalContext context)

View file

@ -6,12 +6,18 @@ using Microsoft.Extensions.Logging;
namespace Elsa.Workflows.Core.Middleware.Activities;
/// <summary>
/// An activity execution middleware component that logs information about the activity being executed.
/// </summary>
public class LoggingMiddleware : IActivityExecutionMiddleware
{
private readonly ActivityMiddlewareDelegate _next;
private readonly ILogger _logger;
private readonly Stopwatch _stopwatch;
/// <summary>
/// Constructor.
/// </summary>
public LoggingMiddleware(ActivityMiddlewareDelegate next, ILogger<LoggingMiddleware> logger)
{
_next = next;
@ -19,18 +25,25 @@ public class LoggingMiddleware : IActivityExecutionMiddleware
_stopwatch = new Stopwatch();
}
/// <inheritdoc />
public async ValueTask InvokeAsync(ActivityExecutionContext context)
{
var activity = context.Activity;
_logger.LogDebug("Executing activity {ActivityType}", activity.Type);
_logger.LogInformation("Executing activity {ActivityId}", activity.Id);
_stopwatch.Restart();
await _next(context);
_stopwatch.Stop();
_logger.LogDebug("Executed activity {ActivityType} in {Elapsed}", activity.Type, _stopwatch.Elapsed);
_logger.LogInformation("Executed activity {ActivityId} in {Elapsed}", activity.Id, _stopwatch.Elapsed);
}
}
/// <summary>
/// Extends <see cref="IActivityExecutionPipelineBuilder"/> to install the <see cref="LoggingMiddleware"/> component.
/// </summary>
public static class LoggingMiddlewareExtensions
{
/// <summary>
/// Installs the <see cref="LoggingMiddleware"/> component.
/// </summary>
public static IActivityExecutionPipelineBuilder UseLogging(this IActivityExecutionPipelineBuilder pipelineBuilder) => pipelineBuilder.UseMiddleware<LoggingMiddleware>();
}

View file

@ -266,7 +266,7 @@ public class WorkflowExecutionContext
/// Returns the <see cref="IActivity"/> with the specified ID from the workflow graph.
/// </summary>
public IActivity FindActivityByNodeId(string nodeId) => FindNodeById(nodeId).Activity;
/// <summary>
/// Returns a custom property with the specified key from the <see cref="Properties"/> dictionary.
/// </summary>

View file

@ -46,7 +46,7 @@ internal class DispatchWorkflowRequestHandler :
public async Task<Unit> HandleAsync(DispatchResumeWorkflowsCommand command, CancellationToken cancellationToken)
{
var options = new ResumeWorkflowRuntimeOptions(CorrelationId: command.CorrelationId, Input: command.Input);
var options = new TriggerWorkflowsRuntimeOptions(CorrelationId: command.CorrelationId, Input: command.Input);
await _workflowRuntime.ResumeWorkflowsAsync(command.ActivityTypeName, command.BookmarkPayload, options, cancellationToken);
return Unit.Instance;

View file

@ -4,6 +4,7 @@ using Elsa.Workflows.Core.Models;
using Elsa.Workflows.Core.Notifications;
using Elsa.Workflows.Runtime.Activities;
using Elsa.Workflows.Runtime.Bookmarks;
using Elsa.Workflows.Runtime.Models;
using Elsa.Workflows.Runtime.Services;
using JetBrains.Annotations;
@ -15,9 +16,9 @@ namespace Elsa.Workflows.Runtime.Handlers;
[PublicAPI]
internal class ResumeDispatchWorkflowActivityHandler : INotificationHandler<WorkflowExecuted>
{
private readonly IWorkflowRuntime _workflowRuntime;
private readonly IWorkflowDispatcher _workflowRuntime;
public ResumeDispatchWorkflowActivityHandler(IWorkflowRuntime workflowRuntime)
public ResumeDispatchWorkflowActivityHandler(IWorkflowDispatcher workflowRuntime)
{
_workflowRuntime = workflowRuntime;
}
@ -29,6 +30,7 @@ internal class ResumeDispatchWorkflowActivityHandler : INotificationHandler<Work
var bookmark = new DispatchWorkflowBookmark(notification.WorkflowState.Id);
var activityTypeName = ActivityTypeNameHelper.GenerateTypeName<DispatchWorkflow>();
await _workflowRuntime.TriggerWorkflowsAsync(activityTypeName, bookmark, new TriggerWorkflowsRuntimeOptions(), cancellationToken);
var request = new DispatchResumeWorkflowsRequest(activityTypeName, bookmark);
await _workflowRuntime.DispatchAsync(request, cancellationToken);
}
}

View file

@ -68,6 +68,40 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
return new WorkflowExecutionResult(workflowState.Id, workflowState.Bookmarks);
}
/// <inheritdoc />
public async Task<ICollection<WorkflowExecutionResult>> StartWorkflowsAsync(
string activityTypeName,
object bookmarkPayload,
TriggerWorkflowsRuntimeOptions options,
CancellationToken cancellationToken = default)
{
var results = new List<WorkflowExecutionResult>();
var hash = _hasher.Hash(activityTypeName, bookmarkPayload);
// Start new workflows. Notice that this happens in a process-synchronized fashion to avoid multiple instances from being created.
var sharedResource = $"{nameof(DefaultWorkflowRuntime)}__StartTriggeredWorkflows__{hash}";
await using (await _distributedLockProvider.AcquireLockAsync(sharedResource, TimeSpan.FromMinutes(10), cancellationToken))
{
var triggers = await _triggerStore.FindAsync(hash, cancellationToken);
foreach (var trigger in triggers)
{
var definitionId = trigger.WorkflowDefinitionId;
var startOptions = new StartWorkflowRuntimeOptions(options.CorrelationId, options.Input, VersionOptions.Published, trigger.ActivityId);
var canStartResult = await CanStartWorkflowAsync(definitionId, startOptions, cancellationToken);
// If we can't start the workflow, don't try it.
if (!canStartResult.CanStart)
continue;
var startResult = await StartWorkflowAsync(definitionId, startOptions, cancellationToken);
results.Add(new WorkflowExecutionResult(startResult.InstanceId, startResult.Bookmarks));
}
}
return results;
}
/// <inheritdoc />
public async Task<ResumeWorkflowResult> ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
{
@ -104,52 +138,21 @@ public class DefaultWorkflowRuntime : IWorkflowRuntime
}
/// <inheritdoc />
public async Task<ICollection<WorkflowExecutionResult>> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default)
public async Task<ICollection<WorkflowExecutionResult>> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default)
{
var hash = _hasher.Hash(activityTypeName, bookmarkPayload);
var correlationId = options.CorrelationId;
var bookmarks = correlationId == null ? await _bookmarkStore.FindByHashAsync(hash, cancellationToken) : await _bookmarkStore.FindByCorrelationAndHashAsync(correlationId, hash, cancellationToken);
return await ResumeWorkflowsAsync(bookmarks, options, cancellationToken);
var bookmarks = string.IsNullOrWhiteSpace(correlationId) ? await _bookmarkStore.FindByHashAsync(hash, cancellationToken) : await _bookmarkStore.FindByCorrelationAndHashAsync(correlationId, hash, cancellationToken);
return await ResumeWorkflowsAsync(bookmarks, new ResumeWorkflowRuntimeOptions(options.CorrelationId, Input: options.Input), cancellationToken);
}
/// <inheritdoc />
public async Task<TriggerWorkflowsResult> TriggerWorkflowsAsync(
string activityTypeName,
object bookmarkPayload,
TriggerWorkflowsRuntimeOptions options,
CancellationToken cancellationToken = default)
public async Task<TriggerWorkflowsResult> TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default)
{
var triggeredWorkflows = new List<WorkflowExecutionResult>();
var hash = _hasher.Hash(activityTypeName, bookmarkPayload);
// Start new workflows. Notice that this happens in a process-synchronized fashion to avoid multiple instances being created.
const string sharedResource = $"{nameof(DefaultWorkflowRuntime)}__StartTriggeredWorkflows";
await using (await _distributedLockProvider.AcquireLockAsync(sharedResource, TimeSpan.FromMinutes(10), cancellationToken))
{
var triggers = await _triggerStore.FindAsync(hash, cancellationToken);
foreach (var trigger in triggers)
{
var definitionId = trigger.WorkflowDefinitionId;
var startOptions = new StartWorkflowRuntimeOptions(options.CorrelationId, options.Input, VersionOptions.Published, trigger.ActivityId);
var canStartResult = await CanStartWorkflowAsync(definitionId, startOptions, cancellationToken);
// If we can't start the workflow, don't try it.
if (!canStartResult.CanStart)
continue;
var startResult = await StartWorkflowAsync(definitionId, startOptions, cancellationToken);
triggeredWorkflows.Add(new WorkflowExecutionResult(startResult.InstanceId, startResult.Bookmarks));
}
}
// Resume bookmarks.
var correlationId = options.CorrelationId;
var bookmarks = (string.IsNullOrEmpty(correlationId) ? await _bookmarkStore.FindByHashAsync(hash, cancellationToken) : await _bookmarkStore.FindByCorrelationAndHashAsync(correlationId, hash, cancellationToken)).ToList();
var resumedWorkflows = await ResumeWorkflowsAsync(bookmarks, new ResumeWorkflowRuntimeOptions(options.CorrelationId, Input: options.Input), cancellationToken);
triggeredWorkflows.AddRange(resumedWorkflows.Select(x => new WorkflowExecutionResult(x.InstanceId, x.Bookmarks)));
return new TriggerWorkflowsResult(triggeredWorkflows);
var startedWorkflows = await StartWorkflowsAsync(activityTypeName, bookmarkPayload, options, cancellationToken);
var resumedWorkflows = await ResumeWorkflowsAsync(activityTypeName, bookmarkPayload, options, cancellationToken);
var results = startedWorkflows.Concat(resumedWorkflows).ToList();
return new TriggerWorkflowsResult(results);
}
/// <inheritdoc />

View file

@ -22,6 +22,16 @@ public interface IWorkflowRuntime
/// <param name="options"></param>
/// <param name="cancellationToken"></param>
Task<WorkflowExecutionResult> StartWorkflowAsync(string definitionId, StartWorkflowRuntimeOptions options, CancellationToken cancellationToken = default);
/// <summary>
/// Starts all workflows with triggers matching the specified activity type and bookmark payload.
/// </summary>
/// <returns></returns>
Task<ICollection<WorkflowExecutionResult>> StartWorkflowsAsync(
string activityTypeName,
object bookmarkPayload,
TriggerWorkflowsRuntimeOptions options,
CancellationToken cancellationToken = default);
/// <summary>
/// Resumes an existing workflow instance.
@ -34,7 +44,7 @@ public interface IWorkflowRuntime
/// <summary>
/// Resumes all workflows that are bookmarked on the specified activity type.
/// </summary>
Task<ICollection<WorkflowExecutionResult>> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, ResumeWorkflowRuntimeOptions options, CancellationToken cancellationToken = default);
Task<ICollection<WorkflowExecutionResult>> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default);
/// <summary>
/// Starts all workflows and resumes existing workflow instances based on the specified activity type and bookmark payload.