From 74a5c3cd87c6551464e1f9bc440d925358aa5398 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 24 Jan 2025 13:58:50 +0100 Subject: [PATCH] Refactor workflow dispatching for clarity and efficiency Refactored the handling of child workflow execution, ensuring better readability and streamlining logic for dispatch operations. Added explicit tracking of dispatched instances and improved comments to support maintainability. --- .../Activities/BulkDispatchWorkflows.cs | 30 +++++++------------ .../Activities/DispatchWorkflow.cs | 11 ++++--- .../Activities/ExecuteWorkflow.cs | 21 ++++++------- 3 files changed, 27 insertions(+), 35 deletions(-) diff --git a/src/modules/Elsa.Workflows.Runtime/Activities/BulkDispatchWorkflows.cs b/src/modules/Elsa.Workflows.Runtime/Activities/BulkDispatchWorkflows.cs index ca1ccd2ae..8074be822 100644 --- a/src/modules/Elsa.Workflows.Runtime/Activities/BulkDispatchWorkflows.cs +++ b/src/modules/Elsa.Workflows.Runtime/Activities/BulkDispatchWorkflows.cs @@ -107,6 +107,13 @@ public class BulkDispatchWorkflows : Activity var items = await context.GetItemSource(Items).ToListAsync(context.CancellationToken); var count = items.Count; + // Dispatch the child workflows. + foreach (var item in items) + await DispatchChildWorkflowAsync(context, item, waitForCompletion); + + // Store the number of dispatched instances for tracking. + context.SetProperty(DispatchedInstancesCountKey, count); + // If we need to wait for the child workflows to complete (if any), create a bookmark. if (waitForCompletion && count > 0) { @@ -122,29 +129,14 @@ public class BulkDispatchWorkflows : Activity IncludeActivityInstanceId = false, AutoBurn = false, }; - - // Create bookmarks first. + context.CreateBookmark(bookmarkOptions); - - // Dispatch workflows afterwards. - await DispatchWorkflowsAsync(); } else { // Otherwise, we can complete immediately. - await DispatchWorkflowsAsync(); await context.CompleteActivityWithOutcomesAsync("Done"); } - - return; - - async Task DispatchWorkflowsAsync() - { - foreach (var item in items) - await DispatchChildWorkflowAsync(context, item, waitForCompletion); - - context.SetProperty(DispatchedInstancesCountKey, count); - } } private async ValueTask DispatchChildWorkflowAsync(ActivityExecutionContext context, object item, bool waitForCompletion) @@ -155,7 +147,7 @@ public class BulkDispatchWorkflows : Activity if (workflowGraph == null) throw new($"No published version of workflow definition with ID {workflowDefinitionId} found."); - + var parentInstanceId = context.WorkflowExecutionContext.Id; var input = Input.GetOrDefault(context) ?? new Dictionary(); var channelName = ChannelName.GetOrDefault(context); @@ -164,8 +156,8 @@ public class BulkDispatchWorkflows : Activity { ["ParentInstanceId"] = parentInstanceId }; - - if(waitForCompletion) + + if (waitForCompletion) properties["WaitForCompletion"] = true; var itemDictionary = new Dictionary diff --git a/src/modules/Elsa.Workflows.Runtime/Activities/DispatchWorkflow.cs b/src/modules/Elsa.Workflows.Runtime/Activities/DispatchWorkflow.cs index 4e579d60c..566c79937 100644 --- a/src/modules/Elsa.Workflows.Runtime/Activities/DispatchWorkflow.cs +++ b/src/modules/Elsa.Workflows.Runtime/Activities/DispatchWorkflow.cs @@ -97,10 +97,10 @@ public class DispatchWorkflow : Activity var workflowDefinitionId = WorkflowDefinitionId.Get(context); var workflowDefinitionService = context.GetRequiredService(); var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionId, VersionOptions.Published, context.CancellationToken); - + if (workflowGraph == null) throw new($"No published version of workflow definition with ID {workflowDefinitionId} found."); - + var input = Input.GetOrDefault(context) ?? new Dictionary(); var channelName = ChannelName.GetOrDefault(context); var parentInstanceId = context.WorkflowExecutionContext.Id; @@ -108,8 +108,9 @@ public class DispatchWorkflow : Activity { ["ParentInstanceId"] = parentInstanceId }; - - if(waitForCompletion) + + // If we need to wait for the child workflow to complete, set the property. This will be used by the ResumeDispatchWorkflowActivity handler. + if (waitForCompletion) properties["WaitForCompletion"] = true; input["ParentInstanceId"] = parentInstanceId; @@ -117,8 +118,6 @@ public class DispatchWorkflow : Activity var correlationId = CorrelationId.GetOrDefault(context); var workflowDispatcher = context.GetRequiredService(); var identityGenerator = context.GetRequiredService(); - - var instanceId = identityGenerator.GenerateId(); var request = new DispatchWorkflowDefinitionRequest(workflowGraph.Workflow.Identity.Id) { diff --git a/src/modules/Elsa.Workflows.Runtime/Activities/ExecuteWorkflow.cs b/src/modules/Elsa.Workflows.Runtime/Activities/ExecuteWorkflow.cs index cba602166..4fda5e880 100644 --- a/src/modules/Elsa.Workflows.Runtime/Activities/ExecuteWorkflow.cs +++ b/src/modules/Elsa.Workflows.Runtime/Activities/ExecuteWorkflow.cs @@ -47,7 +47,7 @@ public class ExecuteWorkflow : Activity /// [Input(Description = "The input to send to the workflow.")] public Input?> Input { get; set; } = null!; - + /// /// True to wait for the child workflow to complete before completing this activity. If not set, the child workflow will be executed until it either completes or goes idle before this activity completes. /// @@ -59,14 +59,14 @@ public class ExecuteWorkflow : Activity { var waitForCompletion = WaitForCompletion.Get(context); var result = await ExecuteWorkflowAsync(context, waitForCompletion); - - if(!waitForCompletion || result.Status == WorkflowStatus.Finished) + + if (!waitForCompletion || result.Status == WorkflowStatus.Finished) { context.SetResult(result); await context.CompleteActivityAsync(); return; } - + // Since the child workflow is still running, we need to wait for it to complete using a bookmark. var bookmarkOptions = new CreateBookmarkArgs { @@ -85,7 +85,7 @@ public class ExecuteWorkflow : Activity if (workflowGraph == null) throw new($"No published version of workflow definition with ID {workflowDefinitionId} found."); - + var parentInstanceId = context.WorkflowExecutionContext.Id; var input = Input.GetOrDefault(context) ?? new Dictionary(); var correlationId = CorrelationId.GetOrDefault(context); @@ -95,12 +95,13 @@ public class ExecuteWorkflow : Activity { ["ParentInstanceId"] = parentInstanceId }; - - if(waitForCompletion) + + // If we need to wait for the child workflow to complete, set the property. This will be used by the ResumeExecuteWorkflowActivity to resume the parent workflow. + if (waitForCompletion) properties["WaitForCompletion"] = true; - + input["ParentInstanceId"] = parentInstanceId; - + var options = new RunWorkflowOptions { ParentWorkflowInstanceId = parentInstanceId, @@ -121,7 +122,7 @@ public class ExecuteWorkflow : Activity return info; } - + private async ValueTask OnChildWorkflowCompletedAsync(ActivityExecutionContext context) { var input = context.WorkflowInput;