elsa-core/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs
Sipke Schoorstra 842cf7c162
[codex] Fix console log metadata and type resolution (#7542)
* Avoid null endpoint DTO metadata in tests

* Enforce console logs hub read permission

* Remove unused console logs hub import

* Support mapped endpoint metadata in auth tests

* Reduce console log capture throughput impact

* Address Copilot console logs review

* Refactor task scheduling to support tenant-level background work and enhance logging functionality.

* Introduce ConsoleStreamHook for stdout/stderr tee and enhance logging validation. Adjust test cases and startup warnings for distributed lock provider usage.

* Refactor console logging pipeline with capture optimization and new ConsoleLogsHost; update tests accordingly.

* Add Ansi SGR parser for console logs and associated unit tests

* Remove ANSI color renderings and parsers; integrate ConsoleLogScopeAccessor for improved logging context with workflow instance ID support.

* Address console logs code quality feedback

* Address PR review feedback

* Preserve console logs extension points

* Stabilize console logs host lifecycle

* Address final automated review comments

* Tighten console log capture shutdown

* Address console log review feedback

* Address follow-up review feedback

* Cover final review feedback

* Avoid recursive console provider initialization

* Guard console host lease shutdown

* Preserve console log scope and provider lifetime

* Correlate console log scope fallback

* Tighten console scope correlation

* Expose host services during provider construction

* Redact ANSI-normalized console lines

* Add OpenTelemetry diagnostics backend foundation

* Add OTLP HTTP ingestion parsing

* Document OpenTelemetry diagnostics setup

* Enforce OpenTelemetry hub permissions

* Remove `ConsoleCaptureTee` and related services and tests

* Add OpenTelemetry HTTP ingestion integration test

* Use pipeline contributors for console log context

* Update CShells package versions to 0.0.24-preview.132

* Add OpenTelemetry ingestion security tests

* Add OpenTelemetry API authorization tests

* Filter live console logs by workflow instance

* Add OpenTelemetry hub tests

* Add OpenTelemetry gRPC metadata hook

* Assert OpenTelemetry workflow tags survive ingestion

* Mark OpenTelemetry core build verified

* Enhance console logging with activity execution metadata and extend test coverage.

* Address console logs stream consumption comment

* Wire OpenTelemetry diagnostics into core sample

* Address Core diagnostics review feedback

* Address Core Copilot follow-up feedback

* Add OpenTelemetry metric instrument names

* Address Core Copilot provider feedback

* Address Core Copilot diagnostics follow-up

* Address Core Copilot live feed feedback

* Address Core Copilot store feedback

* Integrate OpenTelemetry for logging, tracing, and metrics in ModularServer and update launch settings and docker-compose configuration.

* Refactor to replace `ConsoleLogStream.Core` with `ConsoleLogStreaming.Core` across codebase and update `ConsoleStreamHook` installation.

* Add diagnostics OpenTelemetry backend

* Fix OpenTelemetry live hub subscription

* Fix modular OpenTelemetry exporter endpoints

* Add CShells logging configuration in appsettings.json

* Remove obsolete unit tests and helper classes

* Restore default activity exception handling

* Simplify type serialization and alias management

This commit refactors the internal type serialization and alias management system to reduce boilerplate, improve robustness, and simplify the developer experience:

-   Removed numerous explicit `ExpressionOptions` type alias registrations across various modules.
-   Updated `TypeJsonConverter` and polymorphic serialization to reliably handle types using assembly-qualified names when a short alias is not explicitly registered.
-   Streamlined `ExcludeFromHashConverter` to strictly adhere to `ExcludeFromHashAttribute` for hash calculations, removing complex `JsonIgnoreCondition` logic.
-   Eliminated several helper classes (`WorkflowJsonTypeResolver`, `WorkflowTypeValidator`, `IWorkflowTypeRegistry`, `WorkflowFactoryDictionary`, `JavaScriptExceptionTypeAliasRegistrar`, `WorkflowRuntimeTypeAliasRegistrar`) and their associated unit tests, simplifying the codebase.

Additionally, this commit introduces a comprehensive markdown document (`product-website-feature-source.md`) outlining Elsa's core features, Studio capabilities, extension ecosystem, and architectural selling points, intended as source material for the product website.

* Refine type serialization for improved robustness and alias handling

This commit further enhances the type serialization and deserialization mechanisms:

*   Centralizes type resolution and alias management through `IWellKnownTypeRegistry` and `WorkflowJsonTypeResolver`.
*   Prioritizes registered type aliases when serializing type metadata in `PolymorphicObjectConverter`, resulting in more concise JSON output.
*   Enhances deserialization in `PolymorphicObjectConverter` and `VariableMapper` to gracefully handle unknown or non-instantiable types, providing fallbacks and logging warnings.
*   Simplifies `TypeJsonConverter` by delegating complex type resolution logic to the `WorkflowJsonTypeResolver`.
*   Adds `JsonArray` to the well-known type aliases for direct recognition.

* Fix console logs packaging and workflow type resolution

* Fix console log metadata and type resolution

* Address Copilot review feedback

* Enhance type resolution, improve console log handling, and update tests

- Streamlined `WorkflowDictionaryExtensions` for better workflow registration validation.
- Refined `ConsoleLogsAuthorizationTests` with the new `SetJsonRequest` helper to improve test requests handling.
- Updated `OrderDefinition` to ignore JSON serialization for `KeySelector`.
- Enhanced `WorkflowRuntimeFeature` for improved workflow registration and type alias configuration.
- Added tests to ensure `ConsoleLogProvider` metadata filtration in various scenarios.
- Improved type serialization logic in `WorkflowJsonTypeResolver`.
- Updated README to fix references related to diagnostics.
- Optimized `ExcludeFromHashConverter` for property serialization conditions.
- Modified `TriggerIndexer` for streamlined trigger management.
- Tested payload checks in `PublishEventTests`.
- Adjusted `Endpoint` in `ConsoleLogs` for automatic JSON request handling.
- Ensured registration of workflow type aliases in `WorkflowsFeature`.

* Restore CLR workflow registration compatibility

* Align JSON island serialization fixtures

* Add Console Logs Services and Enhance Endpoint Handling

- Introduced `ActivityExecutionsEndpointTests` to validate route exposure.
- Added `ConsoleLogCaptureHostedService` for console log streaming.
- Implemented `ConsoleStreamJsonConverter` for JSON conversion of console streams.
- Developed `ElsaConsoleLogRecentBuffer` to handle recent log buffering.
- Updated `ConsoleLogsAuthorizationTests` with new test cases for stream filter mapping.
- Consolidated console log provider dependencies and registration, including recent buffering.
- Enhanced `ElsaConsoleLogProvider` to use recent buffer for filtering.
- Adjusted `Program.cs` for streamlined logging service setup.

* Enhance type resolution and test coverage; streamline console log integration

- Added `ConsoleStreamHook` for streamlined log streaming.
- Updated `WorkflowJsonTypeResolverTests` to improve type resolution and test new scenarios.
- Simplified type resolution by removing trusted assembly checks.

* Fix CI smoke and package restore failures

* Fix Docker smoke image project paths

* Fix Docker Python runtime packages

* Fix Docker CA smoke teardown

* Refresh Elsa roadmap

* Implement background processors and mediation coordination

- Added `BackgroundCommandProcessor`, `BackgroundJobProcessor`, and `BackgroundNotificationProcessor` classes for handling commands, jobs, and notifications, respectively.
- Introduced `MediatorBackgroundProcessingCoordinator` to coordinate the execution of all background processors.
- Implemented `MediatorBackgroundTask` for wrapping `MediatorBackgroundProcessingCoordinator` in `BackgroundTask`.
- Added unit tests for `MediatorBackgroundTask` to ensure proper start and stop behavior.
- Refactored `BackgroundCommandSenderHostedService` to utilize `BackgroundCommandProcessor`.
- Introduced 'elsa-roadmap-refresh' skill configuration for roadmap updates.

* Address workflow type resolution review feedback

* Address follow-up review feedback

* Restore recent console logs execute path

* Address Copilot follow-up review

* Decouple workflow JSON aliases from expressions

* Fix workflow management unit test setup

* Fix console logs recent endpoint handler shape

* Respect workflow JSON strict type aliases

* Remove unused console log contracts reference

* Address Copilot review feedback

* Address Copilot follow-up comments

* Synchronize ring buffer dropped count

* Address background processor strategy replay
2026-05-30 22:52:01 +02:00

233 lines
10 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.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;
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 DistributedLockingOptions _lockingOptions;
private readonly ILogger _logger;
private readonly WorkflowTriggerEqualityComparer _triggerEqualityComparer;
/// <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,
IOptions<DistributedLockingOptions> lockingOptions,
WorkflowTriggerEqualityComparer triggerEqualityComparer,
ILogger<TriggerIndexer> logger)
{
_activityVisitor = activityVisitor;
_expressionEvaluator = expressionEvaluator;
_identityGenerator = identityGenerator;
_triggerStore = triggerStore;
_activityRegistry = activityRegistry;
_notificationSender = notificationSender;
_serviceProvider = serviceProvider;
_hasher = hasher;
_distributedLockProvider = distributedLockProvider;
_lockingOptions = lockingOptions.Value;
_triggerEqualityComparer = triggerEqualityComparer;
_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 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<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);
}
}