From 8f8f8579a44a09ef2c209b4662bdccb31c178896 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Wed, 18 Dec 2024 15:05:46 +0100 Subject: [PATCH 1/2] Add support for waiting for child workflows in ExecuteWorkflow This update introduces the ability to optionally wait for child workflows to complete before finishing the ExecuteWorkflow activity. A bookmark mechanism is used to resume the activity when the child workflow completes. Additionally, a new handler was added to manage the resumption of these workflows upon completion events. --- .../Activities/ExecuteWorkflow.cs | 52 +++++++++++++++---- .../Bookmarks/ExecuteWorkflowPayload.cs | 9 ++++ .../Features/WorkflowRuntimeFeature.cs | 1 + .../Handlers/ResumeExecuteWorkflowActivity.cs | 40 ++++++++++++++ 4 files changed, 91 insertions(+), 11 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Runtime/Bookmarks/ExecuteWorkflowPayload.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Handlers/ResumeExecuteWorkflowActivity.cs diff --git a/src/modules/Elsa.Workflows.Runtime/Activities/ExecuteWorkflow.cs b/src/modules/Elsa.Workflows.Runtime/Activities/ExecuteWorkflow.cs index 661a64e1d..ffc8bd953 100644 --- a/src/modules/Elsa.Workflows.Runtime/Activities/ExecuteWorkflow.cs +++ b/src/modules/Elsa.Workflows.Runtime/Activities/ExecuteWorkflow.cs @@ -5,6 +5,7 @@ using Elsa.Workflows.Attributes; using Elsa.Workflows.Contracts; using Elsa.Workflows.Management; using Elsa.Workflows.Models; +using Elsa.Workflows.Runtime.Bookmarks; using Elsa.Workflows.Runtime.Contracts; using Elsa.Workflows.Runtime.Parameters; using Elsa.Workflows.UIHints; @@ -20,7 +21,7 @@ namespace Elsa.Workflows.Runtime.Activities; public class ExecuteWorkflow : Activity { /// - public ExecuteWorkflow([CallerFilePath] string? source = default, [CallerLineNumber] int? line = default) : base(source, line) + public ExecuteWorkflow([CallerFilePath] string? source = null, [CallerLineNumber] int? line = null) : base(source, line) { } @@ -32,7 +33,7 @@ public class ExecuteWorkflow : Activity Description = "The definition ID of the workflow to execute.", UIHint = InputUIHints.WorkflowDefinitionPicker )] - public Input WorkflowDefinitionId { get; set; } = default!; + public Input WorkflowDefinitionId { get; set; } = null!; /// /// The correlation ID to associate the workflow with. @@ -41,20 +42,41 @@ public class ExecuteWorkflow : Activity DisplayName = "Correlation ID", Description = "The correlation ID to associate the workflow with." )] - public Input CorrelationId { get; set; } = default!; + public Input CorrelationId { get; set; } = null!; /// /// The input to send to the workflow. /// [Input(Description = "The input to send to the workflow.")] - public Input?> Input { get; set; } = default!; + 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. + /// + [Input(Description = "Wait for the child workflow to complete before completing this activity.")] + public Input WaitForCompletion { get; set; } = null!; /// protected override async ValueTask ExecuteAsync(ActivityExecutionContext context) { var result = await ExecuteWorkflowAsync(context); - context.SetResult(result); - await context.CompleteActivityAsync(); + var waitForCompletion = WaitForCompletion.Get(context); + + 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 + { + Callback = OnChildWorkflowCompletedAsync, + Payload = new ExecuteWorkflowPayload(result.WorkflowInstanceId), + IncludeActivityInstanceId = false + }; + context.CreateBookmark(bookmarkOptions); } private async ValueTask ExecuteWorkflowAsync(ActivityExecutionContext context) @@ -70,15 +92,16 @@ public class ExecuteWorkflow : Activity if (workflowGraph == null) throw new Exception($"No published version of workflow definition with ID {workflowDefinitionId} found."); - var options = new StartWorkflowRuntimeParams + var startParams = new StartWorkflowRuntimeParams { - ParentWorkflowInstanceId = context.WorkflowExecutionContext.Id, + InstanceId = identityGenerator.GenerateId(), Input = input, + ParentWorkflowInstanceId = context.WorkflowExecutionContext.Id, + VersionOptions = VersionOptions.SpecificVersion(workflowGraph.Workflow.Identity.Version), CorrelationId = correlationId, - InstanceId = identityGenerator.GenerateId() + CancellationTokens = context.CancellationToken, }; - - var workflowResult = await workflowRuntime.StartWorkflowAsync(workflowDefinitionId, options); + var workflowResult = await workflowRuntime.StartWorkflowAsync(workflowGraph.Workflow.Identity.DefinitionId, startParams); var info = new ExecuteWorkflowResult { WorkflowInstanceId = workflowResult.WorkflowInstanceId, @@ -89,4 +112,11 @@ public class ExecuteWorkflow : Activity return info; } + + private async ValueTask OnChildWorkflowCompletedAsync(ActivityExecutionContext context) + { + var input = context.WorkflowInput; + context.Set(Result, input); + await context.CompleteActivityAsync(); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Bookmarks/ExecuteWorkflowPayload.cs b/src/modules/Elsa.Workflows.Runtime/Bookmarks/ExecuteWorkflowPayload.cs new file mode 100644 index 000000000..c6cbbec68 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Bookmarks/ExecuteWorkflowPayload.cs @@ -0,0 +1,9 @@ +using Elsa.Workflows.Runtime.Activities; + +namespace Elsa.Workflows.Runtime.Bookmarks; + +/// +/// Bookmark payload for the activity. +/// +/// The instance ID of the child workflow that was created by the activity. +public record ExecuteWorkflowPayload(string ChildInstanceId); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs index d84dd8f50..3dc91ed49 100644 --- a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs @@ -275,6 +275,7 @@ public class WorkflowRuntimeFeature : FeatureBase .AddCommandHandler() .AddNotificationHandler() .AddNotificationHandler() + .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() diff --git a/src/modules/Elsa.Workflows.Runtime/Handlers/ResumeExecuteWorkflowActivity.cs b/src/modules/Elsa.Workflows.Runtime/Handlers/ResumeExecuteWorkflowActivity.cs new file mode 100644 index 000000000..3483a7ff3 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Handlers/ResumeExecuteWorkflowActivity.cs @@ -0,0 +1,40 @@ +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Helpers; +using Elsa.Workflows.Notifications; +using Elsa.Workflows.Runtime.Activities; +using Elsa.Workflows.Runtime.Bookmarks; +using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Models; +using JetBrains.Annotations; + +namespace Elsa.Workflows.Runtime.Handlers; + +/// +/// Resumes any blocking activities when its child workflow completes. +/// +[PublicAPI] +internal class ResumeExecuteWorkflowActivity(IWorkflowInbox bookmarkQueue) : INotificationHandler +{ + private static readonly string ActivityTypeName = ActivityTypeNameHelper.GenerateTypeName(); + + public async Task HandleAsync(WorkflowExecuted notification, CancellationToken cancellationToken) + { + var workflowState = notification.WorkflowState; + + if (workflowState.Status != WorkflowStatus.Finished) + return; + + var workflowInstanceId = notification.WorkflowState.Id; + var payload = new ExecuteWorkflowPayload(workflowInstanceId); + var input = workflowState.Output; + + var message = new NewWorkflowInboxMessage + { + ActivityTypeName = ActivityTypeName, + BookmarkPayload = payload, + Input = input, + }; + + await bookmarkQueue.SubmitAsync(message, cancellationToken); + } +} \ No newline at end of file From fb70022e87b492bb70a5cac9d52e2d206d1c82a9 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Wed, 18 Dec 2024 15:37:38 +0100 Subject: [PATCH 2/2] Add `Variable` parameter to `StorageDriverContext` Updated `StorageDriverContext` to include a `Variable` parameter, ensuring more precise context handling for variable-related operations. Adjusted relevant method calls to pass the required `Variable` argument where necessary. --- .../Elsa.Workflows.Core/Contexts/StorageDriverContext.cs | 3 ++- .../Services/VariablePersistenceManager.cs | 6 +++--- 2 files changed, 5 insertions(+), 4 deletions(-) diff --git a/src/modules/Elsa.Workflows.Core/Contexts/StorageDriverContext.cs b/src/modules/Elsa.Workflows.Core/Contexts/StorageDriverContext.cs index ea39f67cf..aed6ed1d5 100644 --- a/src/modules/Elsa.Workflows.Core/Contexts/StorageDriverContext.cs +++ b/src/modules/Elsa.Workflows.Core/Contexts/StorageDriverContext.cs @@ -1,8 +1,9 @@ using Elsa.Workflows.Contracts; +using Elsa.Workflows.Memory; namespace Elsa.Workflows; /// /// Provides context for storage drivers. /// -public record StorageDriverContext(IExecutionContext ExecutionContext, CancellationToken CancellationToken); \ No newline at end of file +public record StorageDriverContext(IExecutionContext ExecutionContext, Variable Variable, CancellationToken CancellationToken); \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Services/VariablePersistenceManager.cs b/src/modules/Elsa.Workflows.Core/Services/VariablePersistenceManager.cs index 0bd89f183..fa9f10b9a 100644 --- a/src/modules/Elsa.Workflows.Core/Services/VariablePersistenceManager.cs +++ b/src/modules/Elsa.Workflows.Core/Services/VariablePersistenceManager.cs @@ -31,7 +31,7 @@ public class VariablePersistenceManager : IVariablePersistenceManager foreach (var variable in variables) { context.ExpressionExecutionContext.Memory.Declare(variable); - var storageDriverContext = new StorageDriverContext(context, cancellationToken); + var storageDriverContext = new StorageDriverContext(context, variable, cancellationToken); var register = context.ExpressionExecutionContext.Memory; var block = EnsureBlock(register, variable); var metadata = (VariableBlockMetadata)block.Metadata!; @@ -62,7 +62,6 @@ public class VariablePersistenceManager : IVariablePersistenceManager foreach (var context in contexts) { var variables = GetLocalVariables(context).ToList(); - var storageDriverContext = new StorageDriverContext(context, cancellationToken); foreach (var variable in variables) { @@ -75,6 +74,7 @@ public class VariablePersistenceManager : IVariablePersistenceManager var id = GetStateId(variable); var value = block.Value; + var storageDriverContext = new StorageDriverContext(context, variable, cancellationToken); if (value == null) await driver.DeleteAsync(id, storageDriverContext); @@ -91,7 +91,6 @@ public class VariablePersistenceManager : IVariablePersistenceManager var register = context.ExpressionExecutionContext.Memory; var variableList = GetLocalVariables(context).ToList(); var cancellationToken = context.CancellationToken; - var storageDriverContext = new StorageDriverContext(context, cancellationToken); foreach (var variable in variableList) { @@ -105,6 +104,7 @@ public class VariablePersistenceManager : IVariablePersistenceManager continue; var id = GetStateId(variable); + var storageDriverContext = new StorageDriverContext(context, variable, cancellationToken); await driver.DeleteAsync(id, storageDriverContext); register.Blocks.Remove(variable.Id); }