< Summary

Information
Class: Elsa.Workflows.Runtime.Activities.BulkDispatchWorkflows
Assembly: Elsa.Workflows.Runtime
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Runtime/Activities/BulkDispatchWorkflows.cs
Line coverage
98%
Covered lines: 122
Uncovered lines: 2
Coverable lines: 124
Total lines: 273
Line coverage: 98.3%
Branch coverage
88%
Covered branches: 32
Total branches: 36
Branch coverage: 88.8%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
get_WorkflowDefinitionId()100%11100%
get_Items()100%11100%
get_DefaultItemInputKey()100%11100%
get_CorrelationIdFunction()100%11100%
get_Input()100%11100%
get_WaitForCompletion()100%11100%
get_StartNewTrace()100%11100%
get_ChannelName()100%11100%
get_ChildCompleted()100%11100%
get_ChildFaulted()100%11100%
ExecuteAsync()100%66100%
DispatchChildWorkflowAsync()85.71%1414100%
OnChildWorkflowCompletedAsync()71.42%141492.59%
OnChildFinishedCompletedAsync()100%11100%
AttemptToCompleteAsync()100%22100%

File(s)

/home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Runtime/Activities/BulkDispatchWorkflows.cs

#LineLine coverage
 1using System.Runtime.CompilerServices;
 2using Elsa.Common.Models;
 3using Elsa.Expressions.Contracts;
 4using Elsa.Expressions.Helpers;
 5using Elsa.Expressions.Models;
 6using Elsa.Extensions;
 7using Elsa.Workflows.Activities.Flowchart.Attributes;
 8using Elsa.Workflows.Attributes;
 9using Elsa.Workflows.Helpers;
 10using Elsa.Workflows.Management;
 11using Elsa.Workflows.Memory;
 12using Elsa.Workflows.Models;
 13using Elsa.Workflows.Options;
 14using Elsa.Workflows.Runtime.Notifications;
 15using Elsa.Workflows.Runtime.Requests;
 16using Elsa.Workflows.Runtime.Stimuli;
 17using Elsa.Workflows.Runtime.UIHints;
 18using Elsa.Workflows.UIHints;
 19using JetBrains.Annotations;
 20
 21namespace Elsa.Workflows.Runtime.Activities;
 22
 23/// <summary>
 24/// Creates new workflow instances of the specified workflow for each item in the data source and dispatches them for ex
 25/// </summary>
 26[Activity("Elsa", "Composition", "Create new workflow instances for each item in the data source and dispatch them for e
 27[FlowNode("Completed", "Canceled", "Done")]
 28[UsedImplicitly]
 29public class BulkDispatchWorkflows : Activity
 30{
 31    private const string DispatchedInstancesCountKey = nameof(DispatchedInstancesCountKey);
 32    private const string CompletedInstancesCountKey = nameof(CompletedInstancesCountKey);
 33
 34    /// <inheritdoc />
 50335    public BulkDispatchWorkflows([CallerFilePath] string? source = null, [CallerLineNumber] int? line = null) : base(sou
 36    {
 50337    }
 38
 39    /// <summary>
 40    /// The definition ID of the workflows to dispatch.
 41    /// </summary>
 42    [Input(
 43        DisplayName = "Workflow Definition",
 44        Description = "The definition ID of the workflows to dispatch.",
 45        UIHint = InputUIHints.WorkflowDefinitionPicker
 46    )]
 187147    public Input<string> WorkflowDefinitionId { get; set; } = null!;
 48
 49    /// <summary>
 50    /// The data source to use for dispatching the workflows.
 51    /// </summary>
 52    [Input(Description = "The data source to use for dispatching the workflows.")]
 186053    public Input<object> Items { get; set; } = null!;
 54
 55    /// <summary>
 56    /// The default key to use for the item input. Will not be used if the Items contain a list of dictionaries.
 57    /// </summary>
 58    [Input(Description = "The default key to use for the input name when sending the current item to the dispatched work
 187059    public Input<string> DefaultItemInputKey { get; set; } = new("Item");
 60
 61    /// <summary>
 62    /// The correlation ID to associate the workflow with.
 63    /// </summary>
 64    [Input(
 65        DisplayName = "Correlation ID Function",
 66        Description = "A function to compute the correlation ID to associate a dispatched workflow with.",
 67        AutoEvaluate = false)]
 145168    public Input<string?>? CorrelationIdFunction { get; set; }
 69
 70    /// <summary>
 71    /// The input to send to the workflows.
 72    /// </summary>
 73    [Input(Description = "Additional input to send to the workflows being dispatched.")]
 137374    public Input<IDictionary<string, object>?> Input { get; set; } = null!;
 75
 76    /// <summary>
 77    /// True to wait for the child workflow to complete before completing this activity, false to "fire and forget".
 78    /// </summary>
 79    [Input(
 80        Description = "Wait for the dispatched workflows to complete before completing this activity.",
 81        DefaultValue = true)]
 236082    public Input<bool> WaitForCompletion { get; set; } = new(true);
 83
 84    /// <summary>
 85    /// Indicates whether a new trace context should be started for the workflow execution.
 86    /// </summary>
 87    [Input(Description = "Start a new trace context when using Open Telemetry.", Category = "Open Telemetry")]
 186088    public Input<bool> StartNewTrace { get; set; } = new(false);
 89
 90    /// <summary>
 91    /// The channel to dispatch the workflow to.
 92    /// </summary>
 93    [Input(
 94        DisplayName = "Channel",
 95        Description = "The channel to dispatch the workflow to.",
 96        UIHint = InputUIHints.DropDown,
 97        UIHandler = typeof(DispatcherChannelOptionsProvider)
 98    )]
 136799    public Input<string?> ChannelName { get; set; } = null!;
 100
 101    /// <summary>
 102    /// An activity to execute when the child workflow finishes.
 103    /// </summary>
 1551104    [Port] public IActivity? ChildCompleted { get; set; }
 105
 106    /// <summary>
 107    /// An activity to execute when the child workflow faults.
 108    /// </summary>
 1596109    [Port] public IActivity? ChildFaulted { get; set; }
 110
 111    /// <inheritdoc />
 112    protected override async ValueTask ExecuteAsync(ActivityExecutionContext context)
 113    {
 9114        var waitForCompletion = WaitForCompletion.GetOrDefault(context);
 9115        var startNewTrace = StartNewTrace.GetOrDefault(context);
 9116        var items = await context.GetItemSource<object>(Items).ToListAsync(context.CancellationToken);
 9117        var count = items.Count;
 118
 119        // Dispatch the child workflows.
 57120        foreach (var item in items)
 20121            await DispatchChildWorkflowAsync(context, item, waitForCompletion, startNewTrace, count);
 122
 123        // Store the number of dispatched instances for tracking.
 8124        context.SetProperty(DispatchedInstancesCountKey, count);
 125
 126        // If we need to wait for the child workflows to complete (if any), create a bookmark.
 8127        if (waitForCompletion && count > 0)
 128        {
 5129            var workflowInstanceId = context.WorkflowExecutionContext.Id;
 5130            var bookmarkOptions = new CreateBookmarkArgs
 5131            {
 5132                Callback = OnChildWorkflowCompletedAsync,
 5133                Stimulus = new BulkDispatchWorkflowsStimulus(workflowInstanceId)
 5134                {
 5135                    ParentInstanceId = context.WorkflowExecutionContext.Id,
 5136                    ScheduledInstanceIdsCount = count
 5137                },
 5138                IncludeActivityInstanceId = false,
 5139                AutoBurn = false,
 5140            };
 141
 5142            context.CreateBookmark(bookmarkOptions);
 143        }
 144        else
 145        {
 146            // Otherwise, we can complete immediately.
 3147            await context.CompleteActivityWithOutcomesAsync("Done");
 148        }
 8149    }
 150
 151    private async ValueTask<string> DispatchChildWorkflowAsync(ActivityExecutionContext context, object item, bool waitF
 152    {
 20153        var workflowDefinitionId = WorkflowDefinitionId.Get(context);
 20154        var workflowDefinitionService = context.GetRequiredService<IWorkflowDefinitionService>();
 20155        var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionId, VersionOptions.
 156
 20157        if (workflowGraph == null)
 1158            throw new($"No published version of workflow definition with ID {workflowDefinitionId} found.");
 159
 19160        var parentInstanceId = context.WorkflowExecutionContext.Id;
 19161        var baseInput = Input.GetOrDefault(context);
 19162        var input = baseInput != null ? new Dictionary<string, object>(baseInput) : new Dictionary<string, object>();
 19163        var channelName = ChannelName.GetOrDefault(context);
 19164        var defaultInputItemKey = DefaultItemInputKey.GetOrDefault(context, () => "Item")!;
 19165        var properties = new Dictionary<string, object>
 19166        {
 19167            ["ParentInstanceId"] = parentInstanceId
 19168        };
 169
 32170        if (waitForCompletion) properties["WaitForCompletion"] = true;
 19171        if (startNewTrace) properties["StartNewTrace"] = true;
 172
 19173        if (waitForCompletion)
 174        {
 13175            var activityTypeName = ActivityTypeNameHelper.GenerateTypeName<BulkDispatchWorkflows>();
 13176            var stimulus = new BulkDispatchWorkflowsStimulus(parentInstanceId)
 13177            {
 13178                ScheduledInstanceIdsCount = scheduledInstancesCount
 13179            };
 13180            var stimulusHash = context.GetRequiredService<IStimulusHasher>().Hash(activityTypeName, stimulus);
 13181            DispatchWorkflowActivationDeniedRoute.Add(properties, activityTypeName, stimulusHash);
 182        }
 183
 19184        var itemDictionary = new Dictionary<string, object>
 19185        {
 19186            [defaultInputItemKey] = item
 19187        };
 188
 19189        var evaluatorOptions = new ExpressionEvaluatorOptions
 19190        {
 19191            Arguments = itemDictionary
 19192        };
 193
 19194        var inputDictionary = item as IDictionary<string, object> ?? itemDictionary;
 19195        input["ParentInstanceId"] = parentInstanceId;
 19196        input.Merge(inputDictionary);
 197
 19198        var workflowDispatcher = context.GetRequiredService<IWorkflowDispatcher>();
 19199        var identityGenerator = context.GetRequiredService<IIdentityGenerator>();
 19200        var evaluator = context.GetRequiredService<IExpressionEvaluator>();
 19201        var correlationId = CorrelationIdFunction != null ? await evaluator.EvaluateAsync<string>(CorrelationIdFunction!
 19202        var instanceId = identityGenerator.GenerateId();
 19203        var request = new DispatchWorkflowDefinitionRequest(workflowGraph.Workflow.Identity.Id)
 19204        {
 19205            ParentWorkflowInstanceId = parentInstanceId,
 19206            Input = input,
 19207            Properties = properties,
 19208            CorrelationId = correlationId,
 19209            InstanceId = instanceId
 19210        };
 19211        var options = new DispatchWorkflowOptions
 19212        {
 19213            Channel = channelName
 19214        };
 215
 19216        await workflowDispatcher.DispatchAsync(request, options, context.CancellationToken);
 19217        return instanceId;
 19218    }
 219
 220    private async ValueTask OnChildWorkflowCompletedAsync(ActivityExecutionContext context)
 221    {
 13222        var input = context.WorkflowInput;
 13223        var workflowInstanceId = input["WorkflowInstanceId"].ConvertTo<string>()!;
 13224        var workflowSubStatus = input["WorkflowSubStatus"].ConvertTo<WorkflowSubStatus>();
 13225        var cannotStart = input.TryGetValue("CannotStart", out var cannotStartValue) && cannotStartValue is true;
 13226        var finishedInstancesCount = context.GetProperty<long>(CompletedInstancesCountKey) + 1;
 227
 13228        context.SetProperty(CompletedInstancesCountKey, finishedInstancesCount);
 229
 13230        var variables = new List<Variable>();
 13231        if (!cannotStart)
 232        {
 12233            variables.Add(new Variable<string>("ChildInstanceId", workflowInstanceId)
 12234            {
 12235                StorageDriverType = typeof(WorkflowInstanceStorageDriver)
 12236            });
 237        }
 238
 13239        var options = new ScheduleWorkOptions
 13240        {
 13241            Input = input,
 13242            Variables = variables,
 13243            CompletionCallback = OnChildFinishedCompletedAsync
 13244        };
 245
 13246        switch (workflowSubStatus)
 247        {
 4248            case WorkflowSubStatus.Faulted when ChildFaulted is not null:
 4249                await context.ScheduleActivityAsync(ChildFaulted, options);
 4250                return;
 9251            case WorkflowSubStatus.Finished when ChildCompleted is not null:
 0252                await context.ScheduleActivityAsync(ChildCompleted, options);
 0253                return;
 254            default:
 9255                await AttemptToCompleteAsync(context);
 256                break;
 257        }
 13258    }
 259
 260    private async ValueTask OnChildFinishedCompletedAsync(ActivityCompletedContext context)
 261    {
 4262        await AttemptToCompleteAsync(context.TargetContext);
 4263    }
 264
 265    private async ValueTask AttemptToCompleteAsync(ActivityExecutionContext context)
 266    {
 13267        var dispatchedInstancesCount = context.GetProperty<long>(DispatchedInstancesCountKey);
 13268        var finishedInstancesCount = context.GetProperty<long>(CompletedInstancesCountKey);
 269
 13270        if (finishedInstancesCount >= dispatchedInstancesCount)
 5271            await context.CompleteActivityWithOutcomesAsync("Completed", "Done");
 13272    }
 273}