Merge pull request #6226 from elsa-workflows/feat/execute-workflow-activity

Add support for waiting for child workflows in ExecuteWorkflow
This commit is contained in:
Sipke Schoorstra 2024-12-18 15:39:02 +01:00 committed by GitHub
commit fae191d466
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
6 changed files with 96 additions and 15 deletions

View file

@ -1,8 +1,9 @@
using Elsa.Workflows.Contracts;
using Elsa.Workflows.Memory;
namespace Elsa.Workflows;
/// <summary>
/// Provides context for storage drivers.
/// </summary>
public record StorageDriverContext(IExecutionContext ExecutionContext, CancellationToken CancellationToken);
public record StorageDriverContext(IExecutionContext ExecutionContext, Variable Variable, CancellationToken CancellationToken);

View file

@ -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);
}

View file

@ -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<ExecuteWorkflowResult>
{
/// <inheritdoc />
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<ExecuteWorkflowResult>
Description = "The definition ID of the workflow to execute.",
UIHint = InputUIHints.WorkflowDefinitionPicker
)]
public Input<string> WorkflowDefinitionId { get; set; } = default!;
public Input<string> WorkflowDefinitionId { get; set; } = null!;
/// <summary>
/// The correlation ID to associate the workflow with.
@ -41,20 +42,41 @@ public class ExecuteWorkflow : Activity<ExecuteWorkflowResult>
DisplayName = "Correlation ID",
Description = "The correlation ID to associate the workflow with."
)]
public Input<string?> CorrelationId { get; set; } = default!;
public Input<string?> CorrelationId { get; set; } = null!;
/// <summary>
/// The input to send to the workflow.
/// </summary>
[Input(Description = "The input to send to the workflow.")]
public Input<IDictionary<string, object>?> Input { get; set; } = default!;
public Input<IDictionary<string, object>?> Input { get; set; } = null!;
/// <summary>
/// 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.
/// </summary>
[Input(Description = "Wait for the child workflow to complete before completing this activity.")]
public Input<bool> WaitForCompletion { get; set; } = null!;
/// <inheritdoc />
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<ExecuteWorkflowResult> ExecuteWorkflowAsync(ActivityExecutionContext context)
@ -70,15 +92,16 @@ public class ExecuteWorkflow : Activity<ExecuteWorkflowResult>
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<ExecuteWorkflowResult>
return info;
}
private async ValueTask OnChildWorkflowCompletedAsync(ActivityExecutionContext context)
{
var input = context.WorkflowInput;
context.Set(Result, input);
await context.CompleteActivityAsync();
}
}

View file

@ -0,0 +1,9 @@
using Elsa.Workflows.Runtime.Activities;
namespace Elsa.Workflows.Runtime.Bookmarks;
/// <summary>
/// Bookmark payload for the <see cref="ExecuteWorkflow"/> activity.
/// </summary>
/// <param name="ChildInstanceId">The instance ID of the child workflow that was created by the <see cref="ExecuteWorkflow"/> activity.</param>
public record ExecuteWorkflowPayload(string ChildInstanceId);

View file

@ -275,6 +275,7 @@ public class WorkflowRuntimeFeature : FeatureBase
.AddCommandHandler<CancelWorkflowsCommandHandler>()
.AddNotificationHandler<ResumeDispatchWorkflowActivity>()
.AddNotificationHandler<ResumeBulkDispatchWorkflowActivity>()
.AddNotificationHandler<ResumeExecuteWorkflowActivity>()
.AddNotificationHandler<IndexTriggers>()
.AddNotificationHandler<CancelBackgroundActivities>()
.AddNotificationHandler<DeleteBookmarks>()

View file

@ -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;
/// <summary>
/// Resumes any blocking <see cref="ExecuteWorkflow"/> activities when its child workflow completes.
/// </summary>
[PublicAPI]
internal class ResumeExecuteWorkflowActivity(IWorkflowInbox bookmarkQueue) : INotificationHandler<WorkflowExecuted>
{
private static readonly string ActivityTypeName = ActivityTypeNameHelper.GenerateTypeName<ExecuteWorkflow>();
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);
}
}