< 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: 125
Uncovered lines: 3
Coverable lines: 128
Total lines: 252
Line coverage: 97.6%
Branch coverage
93%
Covered branches: 30
Total branches: 32
Branch coverage: 93.7%
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()100%22100%
GetTriggersAsync()100%11100%
DeleteTriggersAsync()100%11100%
GetCurrentTriggersAsync()100%11100%
GetTriggersInternalAsync()100%44100%
CreateWorkflowTriggersAsync()91.66%121291.3%
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>
 75143    public TriggerIndexer(
 75144        IActivityVisitor activityVisitor,
 75145        IWorkflowDefinitionService workflowDefinitionService,
 75146        IExpressionEvaluator expressionEvaluator,
 75147        IIdentityGenerator identityGenerator,
 75148        ITriggerStore triggerStore,
 75149        IActivityRegistry activityRegistry,
 75150        INotificationSender notificationSender,
 75151        IServiceProvider serviceProvider,
 75152        IStimulusHasher hasher,
 75153        IDistributedLockProvider distributedLockProvider,
 75154        ISerializationTypeRegistry workflowJsonTypeRegistry,
 75155        IOptions<DistributedLockingOptions> lockingOptions,
 75156        ILogger<TriggerIndexer> logger)
 57    {
 75158        _activityVisitor = activityVisitor;
 75159        _expressionEvaluator = expressionEvaluator;
 75160        _identityGenerator = identityGenerator;
 75161        _triggerStore = triggerStore;
 75162        _activityRegistry = activityRegistry;
 75163        _notificationSender = notificationSender;
 75164        _serviceProvider = serviceProvider;
 75165        _hasher = hasher;
 75166        _distributedLockProvider = distributedLockProvider;
 75167        _triggerEqualityComparer = new WorkflowTriggerEqualityComparer(workflowJsonTypeRegistry);
 75168        _lockingOptions = lockingOptions.Value;
 75169        _logger = logger;
 75170        _workflowDefinitionService = workflowDefinitionService;
 75171    }
 72
 73    /// <inheritdoc />
 74    public async Task DeleteTriggersAsync(TriggerFilter filter, CancellationToken cancellationToken = default)
 75    {
 976        var triggers = (await _triggerStore.FindManyAsync(filter, cancellationToken)).ToList();
 1677        var workflowDefinitionVersionIds = triggers.Select(x => x.WorkflowDefinitionVersionId).Distinct().ToList();
 78
 3279        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        }
 995    }
 96
 97    /// <inheritdoc />
 98    public async Task<IndexedWorkflowTriggers> IndexTriggersAsync(WorkflowDefinition definition, CancellationToken cance
 99    {
 1670100        var workflowGraph = await _workflowDefinitionService.MaterializeWorkflowAsync(definition, cancellationToken);
 1670101        return await IndexTriggersAsync(workflowGraph.Workflow, cancellationToken);
 1670102    }
 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
 1670108        var lockResource = $"trigger-indexer:{workflow.Identity.DefinitionId}";
 1670109        await using (await _distributedLockProvider.AcquireLockAsync(lockResource, _lockingOptions.LockAcquisitionTimeou
 110        {
 1670111            return await IndexTriggersInternalAsync(workflow, cancellationToken);
 112        }
 1670113    }
 114
 115    private async Task<IndexedWorkflowTriggers> IndexTriggersInternalAsync(Workflow workflow, CancellationToken cancella
 116    {
 117        // Get current triggers
 1670118        var currentTriggers = await GetCurrentTriggersAsync(workflow.Identity.DefinitionId, cancellationToken).ToList();
 119
 120        // Collect new triggers **if the workflow is published**.
 1670121        var newTriggers = workflow.Publication.IsPublished
 1670122            ? await GetTriggersInternalAsync(workflow, cancellationToken).ToListAsync(cancellationToken)
 1670123            : new(0);
 124
 125        // Diff triggers.
 1670126        var diff = Diff.For(currentTriggers, newTriggers, _triggerEqualityComparer);
 127
 128        // Replace triggers for the specified workflow.
 1670129        await _triggerStore.ReplaceAsync(diff.Removed, diff.Added, cancellationToken);
 130
 1670131        var indexedWorkflow = new IndexedWorkflowTriggers(workflow, diff.Added, diff.Removed, diff.Unchanged);
 132
 133        // Publish event.
 1670134        await _notificationSender.SendAsync(new WorkflowTriggersIndexed(indexedWorkflow), cancellationToken);
 1670135        return indexedWorkflow;
 1670136    }
 137
 138    /// <inheritdoc />
 139    public async Task<IEnumerable<StoredTrigger>> GetTriggersAsync(Workflow workflow, CancellationToken cancellationToke
 140    {
 76141        return await GetTriggersInternalAsync(workflow, cancellationToken).ToListAsync(cancellationToken);
 76142    }
 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    {
 1673156        var filter = new TriggerFilter
 1673157        {
 1673158            WorkflowDefinitionId = workflowDefinitionId
 1673159        };
 1673160        return await _triggerStore.FindManyAsync(filter, cancellationToken);
 1673161    }
 162
 163    private async IAsyncEnumerable<StoredTrigger> GetTriggersInternalAsync(Workflow workflow, [EnumeratorCancellation] C
 164    {
 1745165        var context = new WorkflowIndexingContext(workflow, cancellationToken);
 1745166        var nodes = await _activityVisitor.VisitAsync(workflow.Root, cancellationToken);
 167
 168        // Get a list of trigger activities that are configured as "startable".
 1745169        var triggerActivities = nodes
 1745170            .Flatten()
 5645171            .Where(x => x.Activity.GetCanStartWorkflow() && x.Activity is ITrigger)
 357172            .Select(x => x.Activity)
 1745173            .Cast<ITrigger>()
 1745174            .ToList();
 175
 176        // For each trigger activity, create a trigger.
 4204177        foreach (var triggerActivity in triggerActivities)
 178        {
 357179            var triggers = await CreateWorkflowTriggersAsync(context, triggerActivity);
 180
 1546181            foreach (var trigger in triggers)
 416182                yield return trigger;
 183        }
 1745184    }
 185
 186    private async Task<ICollection<StoredTrigger>> CreateWorkflowTriggersAsync(WorkflowIndexingContext context, ITrigger
 187    {
 357188        var workflow = context.Workflow;
 357189        var cancellationToken = context.CancellationToken;
 357190        var activityTypeName = trigger.Type;
 357191        var triggerDescriptor = _activityRegistry.Find(activityTypeName, trigger.Version);
 192
 357193        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
 357199        var expressionExecutionContext = await trigger.CreateExpressionExecutionContextAsync(triggerDescriptor, _service
 357200        var triggerIndexingContext = new TriggerIndexingContext(context, expressionExecutionContext, trigger, cancellati
 357201        var triggerData = await TryGetTriggerDataAsync(trigger, triggerIndexingContext);
 357202        var defaultTriggerName = triggerIndexingContext.TriggerName;
 203
 204        // A trigger that completed, returned no payloads, and declared that deliberate gets no row at all. One that thr
 357205        if (triggerData is { Count: 0 } && triggerIndexingContext.RegistersNoTriggers)
 1206            return new List<StoredTrigger>(0);
 207
 356208        triggerData ??= [];
 209
 210        // If no trigger payloads were returned, create a null payload.
 359211        if (!triggerData.Any()) triggerData.Add(null!);
 212
 356213        var triggers = triggerData.Select(payload =>
 356214        {
 356215            // A payload can carry its own stimulus name. If it does not, the trigger's shared name applies.
 356216            // Name and payload are always taken from the same source so that the hash matches the name stored alongside
 416217            var namedPayload = payload as NamedTriggerPayload;
 416218            var triggerName = namedPayload != null ? namedPayload.Name : defaultTriggerName;
 416219            var stimulus = namedPayload != null ? namedPayload.Payload : payload;
 356220
 416221            return new StoredTrigger
 416222            {
 416223                Id = _identityGenerator.GenerateId(),
 416224                WorkflowDefinitionId = workflow.Identity.DefinitionId,
 416225                WorkflowDefinitionVersionId = workflow.Identity.Id,
 416226                Name = triggerName,
 416227                ActivityId = trigger.Id,
 416228                Hash = _hasher.Hash(triggerName, stimulus),
 416229                Payload = stimulus
 416230            };
 356231        });
 232
 356233        return triggers.ToList();
 357234    }
 235
 236    /// <summary>
 237    /// Returns the trigger's payloads, or <c>null</c> when it throws, so that a failure is told apart from a trigger th
 238    /// </summary>
 239    private async Task<List<object>?> TryGetTriggerDataAsync(ITrigger trigger, TriggerIndexingContext context)
 240    {
 241        try
 242        {
 357243            return (await trigger.GetTriggerPayloadsAsync(context)).ToList();
 244        }
 2245        catch (Exception e)
 246        {
 2247            _logger.LogWarning(e, "Failed to get trigger data for activity {ActivityId}", trigger.Id);
 2248        }
 249
 2250        return null;
 357251    }
 252}