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.
This commit is contained in:
parent
fbd16f2d74
commit
74a5c3cd87
|
|
@ -107,6 +107,13 @@ public class BulkDispatchWorkflows : Activity
|
|||
var items = await context.GetItemSource<object>(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<string> 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<string, object>();
|
||||
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<string, object>
|
||||
|
|
|
|||
|
|
@ -97,10 +97,10 @@ public class DispatchWorkflow : Activity<object>
|
|||
var workflowDefinitionId = WorkflowDefinitionId.Get(context);
|
||||
var workflowDefinitionService = context.GetRequiredService<IWorkflowDefinitionService>();
|
||||
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<string, object>();
|
||||
var channelName = ChannelName.GetOrDefault(context);
|
||||
var parentInstanceId = context.WorkflowExecutionContext.Id;
|
||||
|
|
@ -108,8 +108,9 @@ public class DispatchWorkflow : Activity<object>
|
|||
{
|
||||
["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<object>
|
|||
var correlationId = CorrelationId.GetOrDefault(context);
|
||||
var workflowDispatcher = context.GetRequiredService<IWorkflowDispatcher>();
|
||||
var identityGenerator = context.GetRequiredService<IIdentityGenerator>();
|
||||
|
||||
|
||||
var instanceId = identityGenerator.GenerateId();
|
||||
var request = new DispatchWorkflowDefinitionRequest(workflowGraph.Workflow.Identity.Id)
|
||||
{
|
||||
|
|
|
|||
|
|
@ -47,7 +47,7 @@ public class ExecuteWorkflow : Activity<ExecuteWorkflowResult>
|
|||
/// </summary>
|
||||
[Input(Description = "The input to send to the workflow.")]
|
||||
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>
|
||||
|
|
@ -59,14 +59,14 @@ public class ExecuteWorkflow : Activity<ExecuteWorkflowResult>
|
|||
{
|
||||
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<ExecuteWorkflowResult>
|
|||
|
||||
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<string, object>();
|
||||
var correlationId = CorrelationId.GetOrDefault(context);
|
||||
|
|
@ -95,12 +95,13 @@ public class ExecuteWorkflow : Activity<ExecuteWorkflowResult>
|
|||
{
|
||||
["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<ExecuteWorkflowResult>
|
|||
|
||||
return info;
|
||||
}
|
||||
|
||||
|
||||
private async ValueTask OnChildWorkflowCompletedAsync(ActivityExecutionContext context)
|
||||
{
|
||||
var input = context.WorkflowInput;
|
||||
|
|
|
|||
Loading…
Reference in a new issue