< Summary

Information
Class: Elsa.Workflows.Runtime.StimulusSender
Assembly: Elsa.Workflows.Runtime
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Runtime/Services/StimulusSender.cs
Line coverage
97%
Covered lines: 80
Uncovered lines: 2
Coverable lines: 82
Total lines: 136
Line coverage: 97.5%
Branch coverage
97%
Covered branches: 45
Total branches: 46
Branch coverage: 97.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%
SendAsync(...)100%22100%
SendAsync(...)100%11100%
SendAsync()100%88100%
TriggerNewWorkflowsAsync()92.85%141492.59%
ResumeExistingWorkflowsAsync()100%2222100%

File(s)

/home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Runtime/Services/StimulusSender.cs

#LineLine coverage
 1using Elsa.Workflows.Runtime.Filters;
 2using Elsa.Workflows.Runtime.Messages;
 3using Elsa.Workflows.Runtime.Results;
 4using Microsoft.Extensions.Logging;
 5using Open.Linq.AsyncExtensions;
 6
 7namespace Elsa.Workflows.Runtime;
 8
 9/// <inheritdoc />
 74410public class StimulusSender(
 74411    IStimulusHasher stimulusHasher,
 74412    ITriggerBoundWorkflowService triggerBoundWorkflowService,
 74413    IWorkflowResumer workflowResumer,
 74414    IBookmarkQueue bookmarkQueue,
 74415    ITriggerInvoker triggerInvoker,
 74416    ILogger<StimulusSender> logger) : IStimulusSender
 17{
 18    /// <inheritdoc />
 19    public Task<SendStimulusResult> SendAsync(string activityTypeName, object stimulus, StimulusMetadata? metadata = nul
 20    {
 1221        var stimulusHash = stimulusHasher.Hash(activityTypeName, stimulus, metadata?.ActivityInstanceId);
 1222        return SendAsync(stimulusHash, metadata, activityTypeName, cancellationToken);
 23    }
 24
 25    /// <inheritdoc />
 26    public Task<SendStimulusResult> SendAsync(string stimulusHash, StimulusMetadata? metadata = null, CancellationToken 
 27    {
 428        return SendAsync(stimulusHash, metadata, activityTypeName: null, cancellationToken);
 29    }
 30
 31    private async Task<SendStimulusResult> SendAsync(string stimulusHash, StimulusMetadata? metadata, string? activityTy
 32    {
 1633        var responses = new List<RunWorkflowInstanceResponse>();
 34
 1635        if (metadata == null || (metadata.WorkflowInstanceId == null && metadata.BookmarkId == null && metadata.Activity
 36        {
 637            var triggered = await TriggerNewWorkflowsAsync(stimulusHash, metadata, cancellationToken);
 638            responses.AddRange(triggered);
 39        }
 40
 1641        var resumed = await ResumeExistingWorkflowsAsync(stimulusHash, metadata, activityTypeName, cancellationToken);
 1642        responses.AddRange(resumed);
 1643        return new(responses);
 1644    }
 45
 46    private async Task<ICollection<RunWorkflowInstanceResponse>> TriggerNewWorkflowsAsync(string stimulusHash, StimulusM
 47    {
 648        var triggerBoundWorkflows = await triggerBoundWorkflowService.FindManyAsync(stimulusHash, cancellationToken).ToL
 649        var correlationId = metadata?.CorrelationId;
 650        var input = metadata?.Input;
 651        var properties = metadata?.Properties;
 652        var parentId = metadata?.ParentWorkflowInstanceId;
 653        var responses = new List<RunWorkflowInstanceResponse>();
 54
 2455        foreach (var triggerBoundWorkflow in triggerBoundWorkflows)
 56        {
 657            var workflowGraph = triggerBoundWorkflow.WorkflowGraph;
 658            var workflow = workflowGraph.Workflow;
 59
 2460            foreach (var trigger in triggerBoundWorkflow.Triggers)
 61            {
 662                var triggerRequest = new InvokeTriggerRequest
 663                {
 664                    CorrelationId = correlationId,
 665                    Workflow = workflow,
 666                    ActivityId = trigger.ActivityId,
 667                    Input = input,
 668                    Properties = properties,
 669                    ParentWorkflowInstanceId = parentId
 670                };
 71
 672                var response = await triggerInvoker.InvokeAsync(triggerRequest, cancellationToken);
 73
 674                if (response.CannotStart)
 75                {
 076                    logger.LogWarning("Workflow activation strategy disallowed starting workflow {WorkflowDefinitionId}"
 077                    continue;
 78                }
 79
 680                responses.Add(response.ToRunWorkflowInstanceResponse());
 81            }
 682        }
 83
 684        return responses;
 685    }
 86
 87    private async Task<ICollection<RunWorkflowInstanceResponse>> ResumeExistingWorkflowsAsync(string stimulusHash, Stimu
 88    {
 1689        var input = metadata?.Input;
 1690        var properties = metadata?.Properties;
 91
 1692        var bookmarkFilter = new BookmarkFilter
 1693        {
 1694            Hash = stimulusHash,
 1695            CorrelationId = metadata?.CorrelationId,
 1696            WorkflowInstanceId = metadata?.WorkflowInstanceId,
 1697            ActivityInstanceId = metadata?.ActivityInstanceId,
 1698            BookmarkId = metadata?.BookmarkId
 1699        };
 16100        var responses = (await workflowResumer.ResumeAsync(bookmarkFilter, new()
 16101        {
 16102            Input = input,
 16103            Properties = properties
 16104        }, cancellationToken)).ToList();
 105
 16106        if (responses.Count > 0)
 107        {
 5108            logger.LogDebug("Successfully resumed {WorkflowCount} workflow instances using stimulus {StimulusHash}", res
 5109            return responses;
 110        }
 111
 112        // If no bookmarks were matched, enqueue the request in case a matching bookmark is created in the near future.
 11113        var workflowInstanceId = metadata?.WorkflowInstanceId;
 114
 11115        var bookmarkQueueItem = new NewBookmarkQueueItem
 11116        {
 11117            WorkflowInstanceId = workflowInstanceId,
 11118            BookmarkId = metadata?.BookmarkId,
 11119            CorrelationId = metadata?.CorrelationId,
 11120            StimulusHash = stimulusHash,
 11121            ActivityInstanceId = metadata?.ActivityInstanceId,
 11122            ActivityTypeName = activityTypeName,
 11123            Options = new()
 11124            {
 11125                Input = input,
 11126                Properties = properties
 11127            }
 11128        };
 129
 11130        logger.LogDebug("Bookmark queue item enqueued with stimulus: {StimulusHash}", bookmarkQueueItem.StimulusHash);
 131
 11132        await bookmarkQueue.EnqueueAsync(bookmarkQueueItem, cancellationToken);
 133
 11134        return responses;
 16135    }
 136}