2022-01-04 08:42:12 +00:00
|
|
|
using System.Runtime.CompilerServices;
|
2022-01-14 12:54:17 +00:00
|
|
|
using System.Text.Json;
|
2022-05-26 10:47:31 +00:00
|
|
|
using Elsa.Expressions.Models;
|
|
|
|
|
using Elsa.Expressions.Services;
|
2022-04-25 09:43:23 +00:00
|
|
|
using Elsa.Mediator.Services;
|
2022-05-26 10:47:31 +00:00
|
|
|
using Elsa.Workflows.Core;
|
|
|
|
|
using Elsa.Workflows.Core.Helpers;
|
|
|
|
|
using Elsa.Workflows.Core.Models;
|
|
|
|
|
using Elsa.Workflows.Core.Services;
|
|
|
|
|
using Elsa.Workflows.Persistence.Comparers;
|
|
|
|
|
using Elsa.Workflows.Persistence.Entities;
|
|
|
|
|
using Elsa.Workflows.Persistence.Services;
|
|
|
|
|
using Elsa.Workflows.Runtime.Models;
|
|
|
|
|
using Elsa.Workflows.Runtime.Services;
|
2022-01-04 08:42:12 +00:00
|
|
|
using Microsoft.Extensions.Logging;
|
2022-04-30 13:10:11 +00:00
|
|
|
using Open.Linq.AsyncExtensions;
|
2022-01-04 08:42:12 +00:00
|
|
|
|
2022-05-26 10:47:31 +00:00
|
|
|
namespace Elsa.Workflows.Runtime.Implementations;
|
2022-01-04 08:42:12 +00:00
|
|
|
|
|
|
|
|
public class TriggerIndexer : ITriggerIndexer
|
|
|
|
|
{
|
2022-03-04 20:42:27 +00:00
|
|
|
private readonly IActivityWalker _activityWalker;
|
2022-04-30 09:10:08 +00:00
|
|
|
private readonly IWorkflowDefinitionService _workflowDefinitionService;
|
2022-01-04 08:42:12 +00:00
|
|
|
private readonly IExpressionEvaluator _expressionEvaluator;
|
2022-01-14 12:54:17 +00:00
|
|
|
private readonly IIdentityGenerator _identityGenerator;
|
2022-04-30 13:10:11 +00:00
|
|
|
private readonly IWorkflowTriggerStore _workflowTriggerStore;
|
2022-01-14 12:54:17 +00:00
|
|
|
private readonly IEventPublisher _eventPublisher;
|
2022-01-13 21:38:01 +00:00
|
|
|
private readonly IServiceProvider _serviceProvider;
|
2022-01-04 08:42:12 +00:00
|
|
|
private readonly IHasher _hasher;
|
|
|
|
|
private readonly ILogger _logger;
|
|
|
|
|
|
|
|
|
|
public TriggerIndexer(
|
2022-03-04 20:42:27 +00:00
|
|
|
IActivityWalker activityWalker,
|
2022-04-30 09:10:08 +00:00
|
|
|
IWorkflowDefinitionService workflowDefinitionService,
|
2022-01-04 08:42:12 +00:00
|
|
|
IExpressionEvaluator expressionEvaluator,
|
2022-01-14 12:54:17 +00:00
|
|
|
IIdentityGenerator identityGenerator,
|
2022-04-30 13:10:11 +00:00
|
|
|
IWorkflowTriggerStore workflowTriggerStore,
|
2022-01-14 12:54:17 +00:00
|
|
|
IEventPublisher eventPublisher,
|
2022-01-13 21:38:01 +00:00
|
|
|
IServiceProvider serviceProvider,
|
2022-01-04 08:42:12 +00:00
|
|
|
IHasher hasher,
|
|
|
|
|
ILogger<TriggerIndexer> logger)
|
|
|
|
|
{
|
2022-03-04 20:42:27 +00:00
|
|
|
_activityWalker = activityWalker;
|
2022-01-04 08:42:12 +00:00
|
|
|
_expressionEvaluator = expressionEvaluator;
|
2022-01-14 12:54:17 +00:00
|
|
|
_identityGenerator = identityGenerator;
|
2022-04-30 13:10:11 +00:00
|
|
|
_workflowTriggerStore = workflowTriggerStore;
|
2022-01-14 12:54:17 +00:00
|
|
|
_eventPublisher = eventPublisher;
|
2022-01-13 21:38:01 +00:00
|
|
|
_serviceProvider = serviceProvider;
|
2022-01-04 08:42:12 +00:00
|
|
|
_hasher = hasher;
|
|
|
|
|
_logger = logger;
|
2022-04-30 09:10:08 +00:00
|
|
|
_workflowDefinitionService = workflowDefinitionService;
|
2022-01-04 08:42:12 +00:00
|
|
|
}
|
|
|
|
|
|
2022-04-30 09:10:08 +00:00
|
|
|
public async Task<IndexedWorkflowTriggers> IndexTriggersAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default)
|
2022-01-04 08:42:12 +00:00
|
|
|
{
|
2022-04-30 09:10:08 +00:00
|
|
|
var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(definition, cancellationToken);
|
|
|
|
|
return await IndexTriggersAsync(workflow, cancellationToken);
|
2022-01-04 08:42:12 +00:00
|
|
|
}
|
|
|
|
|
|
2022-03-04 11:09:08 +00:00
|
|
|
public async Task<IndexedWorkflowTriggers> IndexTriggersAsync(Workflow workflow, CancellationToken cancellationToken = default)
|
2022-01-04 08:42:12 +00:00
|
|
|
{
|
2022-02-22 11:51:59 +00:00
|
|
|
// Get current triggers
|
2022-04-30 13:10:11 +00:00
|
|
|
var currentTriggers = await GetCurrentTriggersAsync(workflow.Identity.DefinitionId, cancellationToken).ToList();
|
2022-02-22 11:51:59 +00:00
|
|
|
|
|
|
|
|
// Collect new triggers **if workflow is published**.
|
|
|
|
|
var newTriggers = workflow.Publication.IsPublished
|
|
|
|
|
? await GetTriggersAsync(workflow, cancellationToken).ToListAsync(cancellationToken)
|
|
|
|
|
: new List<WorkflowTrigger>(0);
|
|
|
|
|
|
|
|
|
|
// Diff triggers.
|
2022-04-07 09:22:44 +00:00
|
|
|
var diff = Diff.For(currentTriggers, newTriggers, new WorkflowTriggerHashEqualityComparer());
|
2022-01-04 08:42:12 +00:00
|
|
|
|
|
|
|
|
// Replace triggers for the specified workflow.
|
2022-04-30 13:10:11 +00:00
|
|
|
await _workflowTriggerStore.ReplaceAsync(diff.Removed, diff.Added, cancellationToken);
|
2022-02-22 11:51:59 +00:00
|
|
|
|
2022-04-07 09:22:44 +00:00
|
|
|
var indexedWorkflow = new IndexedWorkflowTriggers(workflow, diff.Added, diff.Removed, diff.Unchanged);
|
2022-01-17 21:44:41 +00:00
|
|
|
|
|
|
|
|
// Publish event.
|
2022-02-22 11:51:59 +00:00
|
|
|
await _eventPublisher.PublishAsync(new WorkflowTriggersIndexed(indexedWorkflow), cancellationToken);
|
|
|
|
|
return indexedWorkflow;
|
2022-01-04 08:42:12 +00:00
|
|
|
}
|
|
|
|
|
|
2022-04-30 13:10:11 +00:00
|
|
|
private async Task<IEnumerable<WorkflowTrigger>> GetCurrentTriggersAsync(string workflowDefinitionId, CancellationToken cancellationToken) =>
|
|
|
|
|
await _workflowTriggerStore.FindManyByWorkflowDefinitionIdAsync(workflowDefinitionId, cancellationToken);
|
2022-02-22 11:51:59 +00:00
|
|
|
|
2022-01-04 08:42:12 +00:00
|
|
|
private async IAsyncEnumerable<WorkflowTrigger> GetTriggersAsync(Workflow workflow, [EnumeratorCancellation] CancellationToken cancellationToken = default)
|
|
|
|
|
{
|
2022-03-07 11:01:27 +00:00
|
|
|
var context = new WorkflowIndexingContext(workflow, cancellationToken);
|
2022-03-08 13:40:47 +00:00
|
|
|
|
2022-03-07 11:01:27 +00:00
|
|
|
// Get a list of activities that are configured as "startable".
|
|
|
|
|
var startableNodes = _activityWalker
|
|
|
|
|
.Walk(workflow.Root)
|
|
|
|
|
.Flatten()
|
|
|
|
|
.Where(x => x.Activity.CanStartWorkflow)
|
|
|
|
|
.ToList();
|
|
|
|
|
|
|
|
|
|
// For each startable node, create triggers.
|
|
|
|
|
foreach (var node in startableNodes)
|
2022-01-04 08:42:12 +00:00
|
|
|
{
|
2022-03-07 14:11:03 +00:00
|
|
|
var triggers = await GetTriggersAsync(context, node.Activity);
|
2022-01-04 08:42:12 +00:00
|
|
|
|
|
|
|
|
foreach (var trigger in triggers)
|
|
|
|
|
yield return trigger;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2022-03-07 14:11:03 +00:00
|
|
|
private async Task<IEnumerable<WorkflowTrigger>> GetTriggersAsync(WorkflowIndexingContext context, IActivity activity)
|
2022-03-07 11:01:27 +00:00
|
|
|
{
|
|
|
|
|
// If the activity implements ITrigger, request its trigger data. Otherwise, create one trigger datum.
|
|
|
|
|
if (activity is ITrigger trigger)
|
|
|
|
|
return await CreateWorkflowTriggersAsync(context, trigger);
|
|
|
|
|
|
|
|
|
|
// Else, create a single workflow trigger with no additional data.
|
|
|
|
|
var simpleTrigger = CreateWorkflowTrigger(context, activity);
|
|
|
|
|
|
|
|
|
|
return new[] { simpleTrigger };
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private WorkflowTrigger CreateWorkflowTrigger(WorkflowIndexingContext context, IActivity activity)
|
|
|
|
|
{
|
|
|
|
|
var workflow = context.Workflow;
|
|
|
|
|
return new WorkflowTrigger
|
|
|
|
|
{
|
|
|
|
|
Id = _identityGenerator.GenerateId(),
|
|
|
|
|
WorkflowDefinitionId = workflow.Identity.DefinitionId,
|
|
|
|
|
Name = activity.TypeName
|
|
|
|
|
};
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private async Task<ICollection<WorkflowTrigger>> CreateWorkflowTriggersAsync(WorkflowIndexingContext context, ITrigger trigger)
|
|
|
|
|
{
|
|
|
|
|
var workflow = context.Workflow;
|
|
|
|
|
var cancellationToken = context.CancellationToken;
|
|
|
|
|
var expressionExecutionContext = await CreateExpressionExecutionContextAsync(context, trigger);
|
|
|
|
|
|
|
|
|
|
var triggerIndexingContext = new TriggerIndexingContext(context, expressionExecutionContext, trigger, cancellationToken);
|
|
|
|
|
var triggerData = await TryGetTriggerDataAsync(trigger, triggerIndexingContext);
|
|
|
|
|
var triggerTypeName = trigger.TypeName;
|
|
|
|
|
|
|
|
|
|
var triggers = triggerData.Select(x => new WorkflowTrigger
|
|
|
|
|
{
|
|
|
|
|
Id = _identityGenerator.GenerateId(),
|
|
|
|
|
WorkflowDefinitionId = workflow.Identity.DefinitionId,
|
|
|
|
|
Name = triggerTypeName,
|
|
|
|
|
Hash = _hasher.Hash(x),
|
|
|
|
|
Data = JsonSerializer.Serialize(x)
|
|
|
|
|
});
|
2022-03-08 13:40:47 +00:00
|
|
|
|
2022-03-07 11:01:27 +00:00
|
|
|
return triggers.ToList();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private async Task<ExpressionExecutionContext> CreateExpressionExecutionContextAsync(WorkflowIndexingContext context, ITrigger trigger)
|
2022-01-04 08:42:12 +00:00
|
|
|
{
|
|
|
|
|
var inputs = trigger.GetInputs();
|
2022-06-04 18:01:20 +00:00
|
|
|
var assignedInputs = inputs.Where(x => x.MemoryReference != null!).ToList();
|
2022-01-04 08:42:12 +00:00
|
|
|
var register = context.GetOrCreateRegister(trigger);
|
2022-03-07 11:01:27 +00:00
|
|
|
var cancellationToken = context.CancellationToken;
|
2022-03-22 13:33:15 +00:00
|
|
|
var expressionInput = new Dictionary<string, object>();
|
2022-05-26 10:47:31 +00:00
|
|
|
var transientProperties = new Dictionary<object, object>();
|
|
|
|
|
var applicationProperties = ExpressionExecutionContextExtensions.CreateApplicationPropertiesFrom(context.Workflow, transientProperties, expressionInput);
|
|
|
|
|
var expressionExecutionContext = new ExpressionExecutionContext(_serviceProvider, register, default, applicationProperties, cancellationToken);
|
2022-01-04 08:42:12 +00:00
|
|
|
|
2022-03-07 11:01:27 +00:00
|
|
|
// Evaluate activity inputs before requesting trigger data.
|
2022-01-04 08:42:12 +00:00
|
|
|
foreach (var input in assignedInputs)
|
|
|
|
|
{
|
2022-06-04 18:01:20 +00:00
|
|
|
var locationReference = input.MemoryReference;
|
2022-01-11 21:04:02 +00:00
|
|
|
|
|
|
|
|
try
|
|
|
|
|
{
|
|
|
|
|
var value = await _expressionEvaluator.EvaluateAsync(input, expressionExecutionContext);
|
|
|
|
|
locationReference.Set(expressionExecutionContext, value);
|
|
|
|
|
}
|
|
|
|
|
catch (Exception e)
|
|
|
|
|
{
|
|
|
|
|
_logger.LogWarning(e, "Failed to evaluate '{@Expression}'", input.Expression);
|
|
|
|
|
}
|
2022-01-04 08:42:12 +00:00
|
|
|
}
|
|
|
|
|
|
2022-03-07 11:01:27 +00:00
|
|
|
return expressionExecutionContext;
|
2022-01-04 08:42:12 +00:00
|
|
|
}
|
2022-01-11 21:04:02 +00:00
|
|
|
|
2022-03-07 11:01:27 +00:00
|
|
|
private async Task<ICollection<object>> TryGetTriggerDataAsync(ITrigger trigger, TriggerIndexingContext context)
|
2022-01-11 21:04:02 +00:00
|
|
|
{
|
|
|
|
|
try
|
|
|
|
|
{
|
2022-03-07 11:01:27 +00:00
|
|
|
return (await trigger.GetTriggerDataAsync(context)).ToList();
|
2022-01-11 21:04:02 +00:00
|
|
|
}
|
|
|
|
|
catch (Exception e)
|
|
|
|
|
{
|
2022-03-07 11:01:27 +00:00
|
|
|
_logger.LogWarning(e, "Failed to get trigger data for activity {ActivityId}", trigger.Id);
|
2022-01-11 21:04:02 +00:00
|
|
|
}
|
|
|
|
|
|
2022-01-13 21:38:01 +00:00
|
|
|
return Array.Empty<object>();
|
2022-01-11 21:04:02 +00:00
|
|
|
}
|
2022-01-04 08:42:12 +00:00
|
|
|
}
|