elsa-core/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs
Sipke Schoorstra f9580e7472
fix(runtime): let a trigger index payloads under per-payload stimulus names (#7950)
* fix(runtime): let a trigger index payloads under per-payload stimulus names

TriggerIndexingContext.TriggerName is a single field read once after all
payloads have been collected, so one ITrigger could only ever register its
payloads under one stimulus name. An implementation that assigned the name
more than once - which the stimulus extension methods do as a side effect -
had the last write applied to every row, and since Hash derives from the same
name, the earlier payloads were stored under a hash no publisher computes.

Adds an additive, opt-in path: a payload returned from GetTriggerPayloadsAsync
may be wrapped in NamedTriggerPayload, which carries the stimulus name for that
payload alone. The indexer takes name and payload from the same source, so
Hash always matches the Name stored beside it, and the wrapper is unwrapped
before storage so payload consumers (validators, the trigger diff comparer,
the scheduler) see the payload the trigger produced.

TriggerName keeps its existing meaning as the default for payloads that do not
carry their own, so every existing ITrigger indexes identically: same Name,
same Hash, same Payload, same row count. The empty-payload placeholder row is
left alone.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

* fix(runtime): refuse a nested trigger payload wrapper

NamedTriggerPayload documented that its Payload is never itself a
wrapper, but nothing enforced it. Reject a NamedTriggerPayload whose
payload is another NamedTriggerPayload at construction time, matching
the existing guard against a blank name.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
2026-08-17 23:43:59 +02:00

244 lines
11 KiB
C#

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.Models;
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;
/// <inheritdoc />
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;
/// <summary>
/// Constructor.
/// </summary>
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<DistributedLockingOptions> lockingOptions,
ILogger<TriggerIndexer> 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;
}
/// <inheritdoc />
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);
}
}
}
/// <inheritdoc />
public async Task<IndexedWorkflowTriggers> IndexTriggersAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default)
{
var workflowGraph = await _workflowDefinitionService.MaterializeWorkflowAsync(definition, cancellationToken);
return await IndexTriggersAsync(workflowGraph.Workflow, cancellationToken);
}
/// <inheritdoc />
public async Task<IndexedWorkflowTriggers> 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<IndexedWorkflowTriggers> 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;
}
/// <inheritdoc />
public async Task<IEnumerable<StoredTrigger>> 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<StoredTrigger>(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<IEnumerable<StoredTrigger>> GetCurrentTriggersAsync(string workflowDefinitionId, CancellationToken cancellationToken)
{
var filter = new TriggerFilter
{
WorkflowDefinitionId = workflowDefinitionId
};
return await _triggerStore.FindManyAsync(filter, cancellationToken);
}
private async IAsyncEnumerable<StoredTrigger> 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<ITrigger>()
.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<ICollection<StoredTrigger>> 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<StoredTrigger>(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 defaultTriggerName = triggerIndexingContext.TriggerName;
// If no trigger payloads were returned, create a null payload.
if (!triggerData.Any()) triggerData.Add(null!);
var triggers = triggerData.Select(payload =>
{
// A payload can carry its own stimulus name. If it does not, the trigger's shared name applies.
// Name and payload are always taken from the same source so that the hash matches the name stored alongside it.
var namedPayload = payload as NamedTriggerPayload;
var triggerName = namedPayload != null ? namedPayload.Name : defaultTriggerName;
var stimulus = namedPayload != null ? namedPayload.Payload : payload;
return new StoredTrigger
{
Id = _identityGenerator.GenerateId(),
WorkflowDefinitionId = workflow.Identity.DefinitionId,
WorkflowDefinitionVersionId = workflow.Identity.Id,
Name = triggerName,
ActivityId = trigger.Id,
Hash = _hasher.Hash(triggerName, stimulus),
Payload = stimulus
};
});
return triggers.ToList();
}
private async Task<List<object>> 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);
}
}