diff --git a/src/activities/Elsa.Activities.Temporal.Common/ActivityResults/ScheduleWorkflowResult.cs b/src/activities/Elsa.Activities.Temporal.Common/ActivityResults/ScheduleWorkflowResult.cs index c0057e360..5b6d58276 100644 --- a/src/activities/Elsa.Activities.Temporal.Common/ActivityResults/ScheduleWorkflowResult.cs +++ b/src/activities/Elsa.Activities.Temporal.Common/ActivityResults/ScheduleWorkflowResult.cs @@ -30,7 +30,7 @@ namespace Elsa.Activities.Temporal.Common.ActivityResults await scheduler.ScheduleWorkflowAsync(null, workflowInstanceId, activityId, tenantId, executeAt, null, cancellationToken); } - activityExecutionContext.WorkflowExecutionContext.RegisterTask(nameof(ScheduleWorkflowResult), ScheduleWorkflowAsync); + activityExecutionContext.WorkflowExecutionContext.RegisterTask(ScheduleWorkflowAsync); } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Temporal.Common/Handlers/ScheduleWorkflows.cs b/src/activities/Elsa.Activities.Temporal.Common/Handlers/ScheduleWorkflows.cs deleted file mode 100644 index 6b011a261..000000000 --- a/src/activities/Elsa.Activities.Temporal.Common/Handlers/ScheduleWorkflows.cs +++ /dev/null @@ -1,16 +0,0 @@ -using System.Threading; -using System.Threading.Tasks; -using Elsa.Activities.Temporal.Common.ActivityResults; -using Elsa.Events; -using MediatR; - -namespace Elsa.Activities.Temporal.Common.Handlers -{ - public class ScheduleWorkflows : INotificationHandler - { - public async Task Handle(WorkflowSuspended notification, CancellationToken cancellationToken) - { - await notification.WorkflowExecutionContext.ExecuteRegisteredTasksAsync(nameof(ScheduleWorkflowResult), cancellationToken); - } - } -} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs b/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs index 253d1711e..8cbc00afa 100644 --- a/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs +++ b/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs @@ -76,9 +76,11 @@ namespace Elsa.Services.Models public object? WorkflowContext { get; set; } /// - /// A collection of tasks to execute post-workflow burst execution. + /// A collection of tasks to execute after the workflow is suspended. + /// This is useful for avoiding race conditions, such as sending a message and then waiting for a message to be received using some MessageReceived activity for example. If the workflow didn't get suspended while that message is received, the workflow would get stuck. + /// By designing activities such as `SendMessage` to only actually send the message post-suspension using the Tasks collection, you can be sure that the workflow gets suspended and blocked on the `MessageReceived` activity before a message reply gets received. /// - public IDictionary>> Tasks { get; set; } = new Dictionary>>(); + public ICollection> Tasks { get; private set; } = new List>(); public string? ContextId { @@ -86,30 +88,15 @@ namespace Elsa.Services.Models set => WorkflowInstance.ContextId = value; } - public void RegisterTask(string groupName, Func task) - { - if (!Tasks.ContainsKey(groupName)) - { - var list = new List> { task }; - Tasks[groupName] = list; - } - else - { - Tasks[groupName].Add(task); - } - } + public void RegisterTask(Func task) => Tasks.Add(task); - public IEnumerable> GetRegisteredTasks(string groupName) => - Tasks.ContainsKey(groupName) ? Tasks[groupName] : Enumerable.Empty>(); - - public async ValueTask ExecuteRegisteredTasksAsync(string groupName, CancellationToken cancellationToken = default) + public async ValueTask ProcessRegisteredTasksAsync(CancellationToken cancellationToken = default) { - var tasks = GetRegisteredTasks(groupName); + var tasks = Tasks.ToList(); + Tasks = new List>(); foreach (var task in tasks) await task(this, cancellationToken); - - Tasks.Remove(groupName); } public async Task RemoveBlockingActivityAsync(BlockingActivity blockingActivity) @@ -123,7 +110,7 @@ namespace Elsa.Services.Models public void SetVariable(string name, object? value) => WorkflowInstance.Variables.Set(name, value); public T? GetVariable() => GetVariable(typeof(T).Name); public T? GetVariable(string name) => WorkflowInstance.Variables.Get(name); - + /// /// Gets a variable from across all scopes, starting with the current scope, going up each scope until the requested variable is found. /// @@ -131,16 +118,16 @@ namespace Elsa.Services.Models public object? GetVariable(string name) { var scopes = WorkflowInstance.Scopes.ToList(); - + var mergedVariables = scopes.Select(x => x.Variables).Aggregate(new Variables(), (current, next) => { var combined = current.Data.MergedWith(next.Data); return new Variables(combined); }); - + return mergedVariables.Get(name); } - + /// /// Gets a workflow variable. /// @@ -155,7 +142,7 @@ namespace Elsa.Services.Models public ActivityScope CurrentScope => WorkflowInstance.Scopes.Peek(); public ActivityScope GetScope(string activityId) => WorkflowInstance.Scopes.First(x => x.ActivityId == activityId); - + public ActivityScope GetNamedScope(string activityName) { var activityBlueprint = GetActivityBlueprintByName(activityName)!; diff --git a/src/core/Elsa.Core/Handlers/ExecuteWorkflowsPostSuspensionTasks.cs b/src/core/Elsa.Core/Handlers/ExecuteWorkflowsPostSuspensionTasks.cs new file mode 100644 index 000000000..8bc10e534 --- /dev/null +++ b/src/core/Elsa.Core/Handlers/ExecuteWorkflowsPostSuspensionTasks.cs @@ -0,0 +1,15 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Events; +using MediatR; + +namespace Elsa.Handlers +{ + public class ExecuteWorkflowsPostSuspensionTasks : INotificationHandler + { + public async Task Handle(WorkflowSuspended notification, CancellationToken cancellationToken) + { + await notification.WorkflowExecutionContext.ProcessRegisteredTasksAsync(cancellationToken); + } + } +} \ No newline at end of file