using System.Runtime.CompilerServices; using Elsa.Expressions.Contracts; using Elsa.Extensions; using Elsa.Mediator.Contracts; using Elsa.Workflows.Activities; using Elsa.Workflows.Contracts; using Elsa.Workflows.Helpers; using Elsa.Workflows.Management.Contracts; using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Runtime.Comparers; using Elsa.Workflows.Runtime.Contracts; using Elsa.Workflows.Runtime.Entities; using Elsa.Workflows.Runtime.Filters; using Elsa.Workflows.Runtime.Notifications; using Microsoft.Extensions.Logging; using Open.Linq.AsyncExtensions; namespace Elsa.Workflows.Runtime.Services; /// 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 INotificationSender _notificationSender; private readonly IServiceProvider _serviceProvider; private readonly IBookmarkHasher _hasher; private readonly ILogger _logger; /// /// Constructor. /// public TriggerIndexer( IActivityVisitor activityVisitor, IWorkflowDefinitionService workflowDefinitionService, IExpressionEvaluator expressionEvaluator, IIdentityGenerator identityGenerator, ITriggerStore triggerStore, INotificationSender notificationSender, IServiceProvider serviceProvider, IBookmarkHasher hasher, ILogger logger) { _activityVisitor = activityVisitor; _expressionEvaluator = expressionEvaluator; _identityGenerator = identityGenerator; _triggerStore = triggerStore; _notificationSender = notificationSender; _serviceProvider = serviceProvider; _hasher = hasher; _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 (string workflowDefinitionVersionId in workflowDefinitionVersionIds) { var workflowDefinition = await _workflowDefinitionService.FindWorkflowDefinitionAsync(workflowDefinitionVersionId, cancellationToken); if (workflowDefinition == null) continue; var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken); await DeleteTriggersAsync(workflow, cancellationToken); } } /// public async Task IndexTriggersAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) { var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(definition, cancellationToken); return await IndexTriggersAsync(workflow, cancellationToken); } /// public async Task IndexTriggersAsync(Workflow workflow, CancellationToken cancellationToken = default) { // Get current triggers var currentTriggers = await GetCurrentTriggersAsync(workflow.Identity.DefinitionId, cancellationToken).ToList(); // Collect new triggers **if workflow is published**. var newTriggers = workflow.Publication.IsPublished ? await GetTriggersInternalAsync(workflow, cancellationToken).ToListAsync(cancellationToken) : new List(0); // Diff triggers. var diff = Diff.For(currentTriggers, newTriggers, new WorkflowTriggerEqualityComparer()); // 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(WorkflowDefinition definition, CancellationToken cancellationToken = default) { var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(definition, cancellationToken); return await GetTriggersAsync(workflow, cancellationToken); } /// public async Task> GetTriggersAsync(Workflow workflow, CancellationToken cancellationToken) { return await GetTriggersInternalAsync(workflow, cancellationToken).ToListAsync(cancellationToken); } /// public 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, new WorkflowTriggerEqualityComparer()); await _triggerStore.ReplaceAsync(diff.Removed, diff.Added, cancellationToken); var indexedWorkflow = new IndexedWorkflowTriggers(workflow, emptyTriggerList, currentTriggers, emptyTriggerList); await _notificationSender.SendAsync(new WorkflowTriggersIndexed(indexedWorkflow), cancellationToken); return indexedWorkflow; } 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 expressionExecutionContext = await trigger.CreateExpressionExecutionContextAsync(_serviceProvider, context, _expressionEvaluator, _logger); var triggerIndexingContext = new TriggerIndexingContext(context, expressionExecutionContext, trigger, cancellationToken); var triggerData = await TryGetTriggerDataAsync(trigger, triggerIndexingContext); var triggerTypeName = trigger.Type; // 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 = triggerTypeName, ActivityId = trigger.Id, Hash = _hasher.Hash(triggerTypeName, 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 List(0); } }