using System.Runtime.CompilerServices; using Elsa.Common.DistributedHosting; using Elsa.Expressions.Contracts; using Elsa.Extensions; using Elsa.Mediator.Contracts; using Elsa.Workflows.Activities; using Elsa.Workflows.Helpers; using Elsa.Workflows.Management; using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Runtime.Comparers; using Elsa.Workflows.Runtime.Entities; using Elsa.Workflows.Runtime.Filters; using Elsa.Workflows.Runtime.Notifications; using Medallion.Threading; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using Open.Linq.AsyncExtensions; using Elsa.Common.Serialization; namespace Elsa.Workflows.Runtime; /// public class TriggerIndexer : ITriggerIndexer { private readonly IActivityVisitor _activityVisitor; private readonly IWorkflowDefinitionService _workflowDefinitionService; private readonly IExpressionEvaluator _expressionEvaluator; private readonly IIdentityGenerator _identityGenerator; private readonly ITriggerStore _triggerStore; private readonly IActivityRegistry _activityRegistry; private readonly INotificationSender _notificationSender; private readonly IServiceProvider _serviceProvider; private readonly IStimulusHasher _hasher; private readonly IDistributedLockProvider _distributedLockProvider; private readonly WorkflowTriggerEqualityComparer _triggerEqualityComparer; private readonly DistributedLockingOptions _lockingOptions; private readonly ILogger _logger; /// /// Constructor. /// public TriggerIndexer( IActivityVisitor activityVisitor, IWorkflowDefinitionService workflowDefinitionService, IExpressionEvaluator expressionEvaluator, IIdentityGenerator identityGenerator, ITriggerStore triggerStore, IActivityRegistry activityRegistry, INotificationSender notificationSender, IServiceProvider serviceProvider, IStimulusHasher hasher, IDistributedLockProvider distributedLockProvider, ISerializationTypeRegistry workflowJsonTypeRegistry, IOptions lockingOptions, ILogger logger) { _activityVisitor = activityVisitor; _expressionEvaluator = expressionEvaluator; _identityGenerator = identityGenerator; _triggerStore = triggerStore; _activityRegistry = activityRegistry; _notificationSender = notificationSender; _serviceProvider = serviceProvider; _hasher = hasher; _distributedLockProvider = distributedLockProvider; _triggerEqualityComparer = new WorkflowTriggerEqualityComparer(workflowJsonTypeRegistry); _lockingOptions = lockingOptions.Value; _logger = logger; _workflowDefinitionService = workflowDefinitionService; } /// public async Task DeleteTriggersAsync(TriggerFilter filter, CancellationToken cancellationToken = default) { var triggers = (await _triggerStore.FindManyAsync(filter, cancellationToken)).ToList(); var workflowDefinitionVersionIds = triggers.Select(x => x.WorkflowDefinitionVersionId).Distinct().ToList(); foreach (var workflowDefinitionVersionId in workflowDefinitionVersionIds) { try { var workflowGraph = await _workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionVersionId, cancellationToken); if (workflowGraph == null) continue; await DeleteTriggersAsync(workflowGraph.Workflow, cancellationToken); } catch (Exception ex) { _logger.LogWarning(ex, "Failed to load workflow graph for workflow definition version {WorkflowDefinitionVersionId}. Skipping trigger deletion for this workflow.", workflowDefinitionVersionId); } } } /// public async Task IndexTriggersAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) { var workflowGraph = await _workflowDefinitionService.MaterializeWorkflowAsync(definition, cancellationToken); return await IndexTriggersAsync(workflowGraph.Workflow, cancellationToken); } /// public async Task IndexTriggersAsync(Workflow workflow, CancellationToken cancellationToken = default) { // Use distributed lock to prevent concurrent trigger indexing race conditions var lockResource = $"trigger-indexer:{workflow.Identity.DefinitionId}"; await using (await _distributedLockProvider.AcquireLockAsync(lockResource, _lockingOptions.LockAcquisitionTimeout, cancellationToken)) { return await IndexTriggersInternalAsync(workflow, cancellationToken); } } private async Task IndexTriggersInternalAsync(Workflow workflow, CancellationToken cancellationToken) { // Get current triggers var currentTriggers = await GetCurrentTriggersAsync(workflow.Identity.DefinitionId, cancellationToken).ToList(); // Collect new triggers **if the workflow is published**. var newTriggers = workflow.Publication.IsPublished ? await GetTriggersInternalAsync(workflow, cancellationToken).ToListAsync(cancellationToken) : new(0); // Diff triggers. var diff = Diff.For(currentTriggers, newTriggers, _triggerEqualityComparer); // Replace triggers for the specified workflow. await _triggerStore.ReplaceAsync(diff.Removed, diff.Added, cancellationToken); var indexedWorkflow = new IndexedWorkflowTriggers(workflow, diff.Added, diff.Removed, diff.Unchanged); // Publish event. await _notificationSender.SendAsync(new WorkflowTriggersIndexed(indexedWorkflow), cancellationToken); return indexedWorkflow; } /// public async Task> GetTriggersAsync(Workflow workflow, CancellationToken cancellationToken) { return await GetTriggersInternalAsync(workflow, cancellationToken).ToListAsync(cancellationToken); } private async Task DeleteTriggersAsync(Workflow workflow, CancellationToken cancellationToken = default) { var emptyTriggerList = new List(0); var currentTriggers = await GetCurrentTriggersAsync(workflow.Identity.DefinitionId, cancellationToken).ToList(); var diff = Diff.For(currentTriggers, emptyTriggerList, _triggerEqualityComparer); await _triggerStore.ReplaceAsync(diff.Removed, diff.Added, cancellationToken); var indexedWorkflow = new IndexedWorkflowTriggers(workflow, emptyTriggerList, currentTriggers, emptyTriggerList); await _notificationSender.SendAsync(new WorkflowTriggersIndexed(indexedWorkflow), cancellationToken); } private async Task> GetCurrentTriggersAsync(string workflowDefinitionId, CancellationToken cancellationToken) { var filter = new TriggerFilter { WorkflowDefinitionId = workflowDefinitionId }; return await _triggerStore.FindManyAsync(filter, cancellationToken); } private async IAsyncEnumerable GetTriggersInternalAsync(Workflow workflow, [EnumeratorCancellation] CancellationToken cancellationToken = default) { var context = new WorkflowIndexingContext(workflow, cancellationToken); var nodes = await _activityVisitor.VisitAsync(workflow.Root, cancellationToken); // Get a list of trigger activities that are configured as "startable". var triggerActivities = nodes .Flatten() .Where(x => x.Activity.GetCanStartWorkflow() && x.Activity is ITrigger) .Select(x => x.Activity) .Cast() .ToList(); // For each trigger activity, create a trigger. foreach (var triggerActivity in triggerActivities) { var triggers = await CreateWorkflowTriggersAsync(context, triggerActivity); foreach (var trigger in triggers) yield return trigger; } } private async Task> CreateWorkflowTriggersAsync(WorkflowIndexingContext context, ITrigger trigger) { var workflow = context.Workflow; var cancellationToken = context.CancellationToken; var activityTypeName = trigger.Type; var triggerDescriptor = _activityRegistry.Find(activityTypeName, trigger.Version); if (triggerDescriptor == null) { _logger.LogWarning("Could not find activity descriptor for activity type {ActivityType}", activityTypeName); return new List(0); } var expressionExecutionContext = await trigger.CreateExpressionExecutionContextAsync(triggerDescriptor, _serviceProvider, context, _expressionEvaluator, _logger); var triggerIndexingContext = new TriggerIndexingContext(context, expressionExecutionContext, trigger, cancellationToken); var triggerData = await TryGetTriggerDataAsync(trigger, triggerIndexingContext); var triggerName = triggerIndexingContext.TriggerName; // If no trigger payloads were returned, create a null payload. if (!triggerData.Any()) triggerData.Add(null!); var triggers = triggerData.Select(payload => new StoredTrigger { Id = _identityGenerator.GenerateId(), WorkflowDefinitionId = workflow.Identity.DefinitionId, WorkflowDefinitionVersionId = workflow.Identity.Id, Name = triggerName, ActivityId = trigger.Id, Hash = _hasher.Hash(triggerName, payload), Payload = payload }); return triggers.ToList(); } private async Task> TryGetTriggerDataAsync(ITrigger trigger, TriggerIndexingContext context) { try { return (await trigger.GetTriggerPayloadsAsync(context)).ToList(); } catch (Exception e) { _logger.LogWarning(e, "Failed to get trigger data for activity {ActivityId}", trigger.Id); } return new(0); } }