< Summary

Information
Class: Elsa.Workflows.Runtime.TriggerIndexer
Assembly: Elsa.Workflows.Runtime
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs
Line coverage
97%
Covered lines: 122
Uncovered lines: 3
Coverable lines: 125
Total lines: 243
Line coverage: 97.6%
Branch coverage
87%
Covered branches: 21
Total branches: 24
Branch coverage: 87.5%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
DeleteTriggersAsync()87.5%8891.66%
IndexTriggersAsync()100%11100%
IndexTriggersAsync()100%11100%
IndexTriggersInternalAsync()50%22100%
GetTriggersAsync()100%11100%
DeleteTriggersAsync()100%11100%
GetCurrentTriggersAsync()100%11100%
GetTriggersInternalAsync()100%44100%
CreateWorkflowTriggersAsync()75%4490%
TryGetTriggerDataAsync()100%11100%

File(s)

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

#LineLine coverage
 1using System.Runtime.CompilerServices;
 2using Elsa.Common.DistributedHosting;
 3using Elsa.Expressions.Contracts;
 4using Elsa.Extensions;
 5using Elsa.Mediator.Contracts;
 6using Elsa.Workflows.Activities;
 7using Elsa.Workflows.Helpers;
 8using Elsa.Workflows.Management;
 9using Elsa.Workflows.Management.Entities;
 10using Elsa.Workflows.Models;
 11using Elsa.Workflows.Runtime.Comparers;
 12using Elsa.Workflows.Runtime.Entities;
 13using Elsa.Workflows.Runtime.Filters;
 14using Elsa.Workflows.Runtime.Notifications;
 15using Medallion.Threading;
 16using Microsoft.Extensions.Logging;
 17using Microsoft.Extensions.Options;
 18using Open.Linq.AsyncExtensions;
 19using Elsa.Common.Serialization;
 20
 21namespace Elsa.Workflows.Runtime;
 22
 23/// <inheritdoc />
 24public class TriggerIndexer : ITriggerIndexer
 25{
 26    private readonly IActivityVisitor _activityVisitor;
 27    private readonly IWorkflowDefinitionService _workflowDefinitionService;
 28    private readonly IExpressionEvaluator _expressionEvaluator;
 29    private readonly IIdentityGenerator _identityGenerator;
 30    private readonly ITriggerStore _triggerStore;
 31    private readonly IActivityRegistry _activityRegistry;
 32    private readonly INotificationSender _notificationSender;
 33    private readonly IServiceProvider _serviceProvider;
 34    private readonly IStimulusHasher _hasher;
 35    private readonly IDistributedLockProvider _distributedLockProvider;
 36    private readonly WorkflowTriggerEqualityComparer _triggerEqualityComparer;
 37    private readonly DistributedLockingOptions _lockingOptions;
 38    private readonly ILogger _logger;
 39
 40    /// <summary>
 41    /// Constructor.
 42    /// </summary>
 71243    public TriggerIndexer(
 71244        IActivityVisitor activityVisitor,
 71245        IWorkflowDefinitionService workflowDefinitionService,
 71246        IExpressionEvaluator expressionEvaluator,
 71247        IIdentityGenerator identityGenerator,
 71248        ITriggerStore triggerStore,
 71249        IActivityRegistry activityRegistry,
 71250        INotificationSender notificationSender,
 71251        IServiceProvider serviceProvider,
 71252        IStimulusHasher hasher,
 71253        IDistributedLockProvider distributedLockProvider,
 71254        ISerializationTypeRegistry workflowJsonTypeRegistry,
 71255        IOptions<DistributedLockingOptions> lockingOptions,
 71256        ILogger<TriggerIndexer> logger)
 57    {
 71258        _activityVisitor = activityVisitor;
 71259        _expressionEvaluator = expressionEvaluator;
 71260        _identityGenerator = identityGenerator;
 71261        _triggerStore = triggerStore;
 71262        _activityRegistry = activityRegistry;
 71263        _notificationSender = notificationSender;
 71264        _serviceProvider = serviceProvider;
 71265        _hasher = hasher;
 71266        _distributedLockProvider = distributedLockProvider;
 71267        _triggerEqualityComparer = new WorkflowTriggerEqualityComparer(workflowJsonTypeRegistry);
 71268        _lockingOptions = lockingOptions.Value;
 71269        _logger = logger;
 71270        _workflowDefinitionService = workflowDefinitionService;
 71271    }
 72
 73    /// <inheritdoc />
 74    public async Task DeleteTriggersAsync(TriggerFilter filter, CancellationToken cancellationToken = default)
 75    {
 676        var triggers = (await _triggerStore.FindManyAsync(filter, cancellationToken)).ToList();
 1377        var workflowDefinitionVersionIds = triggers.Select(x => x.WorkflowDefinitionVersionId).Distinct().ToList();
 78
 2679        foreach (var workflowDefinitionVersionId in workflowDefinitionVersionIds)
 80        {
 81            try
 82            {
 783                var workflowGraph = await _workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionVersionId,
 84
 385                if (workflowGraph == null)
 086                    continue;
 87
 388                await DeleteTriggersAsync(workflowGraph.Workflow, cancellationToken);
 389            }
 490            catch (Exception ex)
 91            {
 492                _logger.LogWarning(ex, "Failed to load workflow graph for workflow definition version {WorkflowDefinitio
 493            }
 794        }
 695    }
 96
 97    /// <inheritdoc />
 98    public async Task<IndexedWorkflowTriggers> IndexTriggersAsync(WorkflowDefinition definition, CancellationToken cance
 99    {
 1514100        var workflowGraph = await _workflowDefinitionService.MaterializeWorkflowAsync(definition, cancellationToken);
 1514101        return await IndexTriggersAsync(workflowGraph.Workflow, cancellationToken);
 1514102    }
 103
 104    /// <inheritdoc />
 105    public async Task<IndexedWorkflowTriggers> IndexTriggersAsync(Workflow workflow, CancellationToken cancellationToken
 106    {
 107        // Use distributed lock to prevent concurrent trigger indexing race conditions
 1514108        var lockResource = $"trigger-indexer:{workflow.Identity.DefinitionId}";
 1514109        await using (await _distributedLockProvider.AcquireLockAsync(lockResource, _lockingOptions.LockAcquisitionTimeou
 110        {
 1514111            return await IndexTriggersInternalAsync(workflow, cancellationToken);
 112        }
 1514113    }
 114
 115    private async Task<IndexedWorkflowTriggers> IndexTriggersInternalAsync(Workflow workflow, CancellationToken cancella
 116    {
 117        // Get current triggers
 1514118        var currentTriggers = await GetCurrentTriggersAsync(workflow.Identity.DefinitionId, cancellationToken).ToList();
 119
 120        // Collect new triggers **if the workflow is published**.
 1514121        var newTriggers = workflow.Publication.IsPublished
 1514122            ? await GetTriggersInternalAsync(workflow, cancellationToken).ToListAsync(cancellationToken)
 1514123            : new(0);
 124
 125        // Diff triggers.
 1514126        var diff = Diff.For(currentTriggers, newTriggers, _triggerEqualityComparer);
 127
 128        // Replace triggers for the specified workflow.
 1514129        await _triggerStore.ReplaceAsync(diff.Removed, diff.Added, cancellationToken);
 130
 1514131        var indexedWorkflow = new IndexedWorkflowTriggers(workflow, diff.Added, diff.Removed, diff.Unchanged);
 132
 133        // Publish event.
 1514134        await _notificationSender.SendAsync(new WorkflowTriggersIndexed(indexedWorkflow), cancellationToken);
 1514135        return indexedWorkflow;
 1514136    }
 137
 138    /// <inheritdoc />
 139    public async Task<IEnumerable<StoredTrigger>> GetTriggersAsync(Workflow workflow, CancellationToken cancellationToke
 140    {
 70141        return await GetTriggersInternalAsync(workflow, cancellationToken).ToListAsync(cancellationToken);
 70142    }
 143
 144    private async Task DeleteTriggersAsync(Workflow workflow, CancellationToken cancellationToken = default)
 145    {
 3146        var emptyTriggerList = new List<StoredTrigger>(0);
 3147        var currentTriggers = await GetCurrentTriggersAsync(workflow.Identity.DefinitionId, cancellationToken).ToList();
 3148        var diff = Diff.For(currentTriggers, emptyTriggerList, _triggerEqualityComparer);
 3149        await _triggerStore.ReplaceAsync(diff.Removed, diff.Added, cancellationToken);
 3150        var indexedWorkflow = new IndexedWorkflowTriggers(workflow, emptyTriggerList, currentTriggers, emptyTriggerList)
 3151        await _notificationSender.SendAsync(new WorkflowTriggersIndexed(indexedWorkflow), cancellationToken);
 3152    }
 153
 154    private async Task<IEnumerable<StoredTrigger>> GetCurrentTriggersAsync(string workflowDefinitionId, CancellationToke
 155    {
 1517156        var filter = new TriggerFilter
 1517157        {
 1517158            WorkflowDefinitionId = workflowDefinitionId
 1517159        };
 1517160        return await _triggerStore.FindManyAsync(filter, cancellationToken);
 1517161    }
 162
 163    private async IAsyncEnumerable<StoredTrigger> GetTriggersInternalAsync(Workflow workflow, [EnumeratorCancellation] C
 164    {
 1584165        var context = new WorkflowIndexingContext(workflow, cancellationToken);
 1584166        var nodes = await _activityVisitor.VisitAsync(workflow.Root, cancellationToken);
 167
 168        // Get a list of trigger activities that are configured as "startable".
 1584169        var triggerActivities = nodes
 1584170            .Flatten()
 5375171            .Where(x => x.Activity.GetCanStartWorkflow() && x.Activity is ITrigger)
 352172            .Select(x => x.Activity)
 1584173            .Cast<ITrigger>()
 1584174            .ToList();
 175
 176        // For each trigger activity, create a trigger.
 3872177        foreach (var triggerActivity in triggerActivities)
 178        {
 352179            var triggers = await CreateWorkflowTriggersAsync(context, triggerActivity);
 180
 1528181            foreach (var trigger in triggers)
 412182                yield return trigger;
 183        }
 1584184    }
 185
 186    private async Task<ICollection<StoredTrigger>> CreateWorkflowTriggersAsync(WorkflowIndexingContext context, ITrigger
 187    {
 352188        var workflow = context.Workflow;
 352189        var cancellationToken = context.CancellationToken;
 352190        var activityTypeName = trigger.Type;
 352191        var triggerDescriptor = _activityRegistry.Find(activityTypeName, trigger.Version);
 192
 352193        if (triggerDescriptor == null)
 194        {
 0195            _logger.LogWarning("Could not find activity descriptor for activity type {ActivityType}", activityTypeName);
 0196            return new List<StoredTrigger>(0);
 197        }
 198
 352199        var expressionExecutionContext = await trigger.CreateExpressionExecutionContextAsync(triggerDescriptor, _service
 352200        var triggerIndexingContext = new TriggerIndexingContext(context, expressionExecutionContext, trigger, cancellati
 352201        var triggerData = await TryGetTriggerDataAsync(trigger, triggerIndexingContext);
 352202        var defaultTriggerName = triggerIndexingContext.TriggerName;
 203
 204        // If no trigger payloads were returned, create a null payload.
 354205        if (!triggerData.Any()) triggerData.Add(null!);
 206
 352207        var triggers = triggerData.Select(payload =>
 352208        {
 352209            // A payload can carry its own stimulus name. If it does not, the trigger's shared name applies.
 352210            // Name and payload are always taken from the same source so that the hash matches the name stored alongside
 412211            var namedPayload = payload as NamedTriggerPayload;
 412212            var triggerName = namedPayload != null ? namedPayload.Name : defaultTriggerName;
 412213            var stimulus = namedPayload != null ? namedPayload.Payload : payload;
 352214
 412215            return new StoredTrigger
 412216            {
 412217                Id = _identityGenerator.GenerateId(),
 412218                WorkflowDefinitionId = workflow.Identity.DefinitionId,
 412219                WorkflowDefinitionVersionId = workflow.Identity.Id,
 412220                Name = triggerName,
 412221                ActivityId = trigger.Id,
 412222                Hash = _hasher.Hash(triggerName, stimulus),
 412223                Payload = stimulus
 412224            };
 352225        });
 226
 352227        return triggers.ToList();
 352228    }
 229
 230    private async Task<List<object>> TryGetTriggerDataAsync(ITrigger trigger, TriggerIndexingContext context)
 231    {
 232        try
 233        {
 352234            return (await trigger.GetTriggerPayloadsAsync(context)).ToList();
 235        }
 1236        catch (Exception e)
 237        {
 1238            _logger.LogWarning(e, "Failed to get trigger data for activity {ActivityId}", trigger.Id);
 1239        }
 240
 1241        return new(0);
 352242    }
 243}